How to share a cache in Flink kinesis stream

Viewed 163

I've been using Flink and kinesis analytics recently. I have a stream of data and also I need a cache to be shared with the stream.

To share the cache data with the kinesis stream, it's connected to a broadcast stream. The cache source extends SourceFunction and implements ProcessingTimeCallback. Gets the data from DynamoDB every 300 seconds and broadcast it to the next stream using KeyedBroadcastProcessFuction.

But after adding the broadcast stream (in the previous version I hadn't a cache and I was using KeyedProcessFuction for kinesis stream), when I execute it in kinesis analytics, it keeps restarting about every 1000 seconds without any exception!

I have no configuration with this value and the scenario works fine in between!

Could anybody help me what could be the issue?

1 Answers

My first thought is to wonder if this might be related to checkpointing. Do you have access to the server logs? Flink's logging should make it somewhat clear what's causing the restart.

The reason why I suspect checkpointing is that it occurs at predictable times (and with a long timeout), and using broadcast state can put a lot of pressure on checkpointing. Each parallel instance will checkpoint a full copy of the broadcast state.

Broadcast state has to be kept on-heap, so another possibility is that you are running out of memory.

Related