Real-time data streaming using Wikipedia's RecentChanges API

Viewed 660

I'm lately trying to create a demo on real time streaming using NiFi -> Kafka -> Druid -> Superset. For the purposes of this demo I chose to use Wikipedia's RecentChanges API in order to get asynchronous data of the most recent changes.

I use this URL in order to get a response of changes. I'm calling the API constanlty in order to not miss any changes. This way I get a lot of duplicates that I do not want.

Is there anyway to parameterize this API to fix it for example getting all the changes from the previous second and doing that everysecond or something else to tackle this issue. I'm trying to make a configuration for this uing NiFi, if someone has to add something on that part then visit this discussion on Cloudera.

2 Answers

I want to expand smartse's answer and come up with a solution. You want to put your API request in certain time windows, by shifting the start and end parameters. Windowing might work like this:

  • Initialize start, end timestamp parameters
  • Put those parameters as attributes on the flow
  • Downstream processors can call the API using those parameters
  • After doing that, you have to set start = previous_end + 1 second and end = now

When you determine the new window for the next run, you need the parameters from the previous run. This is why you have to remember those values. You can achieve this using NiFi's distributed map cache.

I've assembled a flow for you:

enter image description here

Zoom into Get next date range:

enter image description here

The end parameter is always now, so you just have to store the start parameter. FetchDistributedMapCache will fetch that for you and put it into stored.state attribute:

enter image description here

Set time range processor will initialize the parameters:

enter image description here

Notice that end is always now and start is either an initial date (for the first run) or the last end parameter plus 1 second. At this point the flow is directed into the Time range output, where you can call your API downstream. Additionally you have to update the stored.value. This happens in the ReplaceText processor:

enter image description here

Finally you update the state:

enter image description here

The lifecycle of the parameters are bound to the cache identifier. When you change the identifier, you start from scratch.

Related