Apache Flink using Java - Performance Issue

Viewed 223

We have a flink application written on Java and running on AWS Kinesis Data Analytics. The application reads the input stream from AWS Managed Service Kafka (kafka topic 1), then apply business logic (some calculations) and finally writes the output to another Kafka topic (kafka topic 2).

The parallelism is 10 and the topic has 15 partitions. The expectation is to process ~20K concurrent data in 5 mins. But after all the optimizations we could bring it to a speed of ~20K concurrent data in 25 mins.

Could you please let me know if there is any other performance optimization that can be implemented to achieve the target.

Is Flink Async I/O going to be an option to further optimize?.

Sample code:-

StreamExecutionEnvironment streamenv =
    StreamExecutionEnvironment.getExecutionEnvironment();

DataStream<ObjectNode> initialStreamData = streamenv
    .addSource(new FlinkKafkaConsumer<>(
        TOPIC_NAME, 
        new ObjectNodeJsonDeSerializerSchema(),
        kafkaConnectProperties);

initialStreamData.print();

DataStream<POJO> rawDataProcess = initialStreamData
    .rebalance()
    .flatMap(new ReProcessingDataProvider())
    .keyBy(value -> value.getPersonId());

rawDataProcess.print();

DataStream<POJO> cgmStream = rawDataProcess
    .keyBy(new ReProcessorKeySelector())
    .rebalance()
    .flatMap(new SgStreamTask());

cgmStream.print();

DataStream<POJO> artfctOverlapStream = null;
artfctOverlapStream = cgmStreamData
    .keyBy(new CGMKeySelector())
    .countWindow(2, 1)
    .apply(new ArtifactOverlapProvider()); //the same person_id key

cgmStreamData.print();

DataStream<POJO> streamWithSgRoc = null;

streamWithSgRoc = artfctOverlapStream
    .keyBy(new CGMKeySelector())
    .countWindow(7, 1)
    .apply(new SgRocProvider()); // the same person_id key 

streamWithSgRoc.print();

DataStream<POJO> cgmExcursionStream = null;

cgmExcursionStream = streamWithSgRoc
    .keyBy(new CGMKeySelector())
    .countWindow(Common.THREE, Common.ONE)
    .apply(new CGMExcursionProviderStream()); //the same person_id key

cgmExcursionStream.print();

cgmExcursionStream
    .addSink(new FlinkKafkaProducer<CGMDataCollector>(
        topicsProperties.getProperty(Common.CGM_EVENT_TOPIC),
        new CGMDataCollectorSchema(),
        kafkaConnectProperties));
1 Answers

I don't see anything in what you've shared that would explain the poor throughput. But since you've asked about async i/o, I wonder if the flatmap is doing some external i/o. If so, that would explain it. If that is the case, then using async i/o should help substantially, provided the external service can handle the increased load.

I'm also wondering why the parallelism is 10, and what resources are available for those 10 slots. Are there enough cores to keep everything moving? With 15 partitions, you have 5 slots handling two partitions each, and the other 5 are handling one partition each. 5, 8, and 15 are more obvious choices for the parallelism. (Of course, if each slot is also hitting an external service in the flatmap, that also needs to be taken into account.)

UPDATED after seeing the code:

There are a couple of easy things you can do to speed this up.

One thing to do is to give the cluster more resources. You can leave the parallelism as it is, but give each slot more cores to work with, by putting the task managers on machines with more cores.

But before doing that take a look at optimizing the pipeline. Both rebalance and keyBy are rather expensive, and you are using them more than is necessary. Having a keyBy followed immediately by a rebalance doesn't make sense, nor does having two keyBy's one immediately after another.

Rebalance does a round-trip re-partitioning of the stream. It is typically done when changing parallelism, or when needed to overcome data skew. Rebalance is almost never used, except implicitly as needed when changing parallelism.

KeyBy does a key-based repartitioning. If one keyBy follows another, the second one undoes whatever the first one did.

Both keyBy and rebalance require serializing and deserializing every event, and sending them through the network stack. This is something you only want to do if it's strictly necessary.

Fixing these rebalance/keyBy issues will reduce the workload on your cluster. If that isn't enough to achieve the desired throughput, giving each slot more cores (so that the various stages of your pipeline can run in parallel) should do the trick.

Related