flink Pipeline not processing Kafka messages after switching source

Viewed 91

I have a use case to build a stateful application. I need to build state with historical data stored in S3 and switch to kafka ( from a particular offset) with continuing historical state.

We have beam pipeline having runner as flink. Below are the steps I am doing. Process all the files and build a state. Stop flink job with savepoint Start flink job with state saved in step2 and switch to kafka as source

My pipeline is not processing any messages in step3. When I check flink UI I do observe the watermark is set as MAX_WATERMARK (9223372036854776000) for stateful transformation block. I am looking for solution if there is way we can override this watermark and set it to required offset.

Below are sample code and topology. For POC I am reading data from local files and I am getting file names from a kafka.

I am using flink version 1.9.3 beam version 2.23.0

try {
    PCollection<String> records = null;
    boolean isSource1 = runningMode.equals("source1") ? true : false;

    if (isSource1) {
        PCollection<String> absolutePaths = pipeline
                .apply("read from source topic 1", KafkaIO.<Long, String>read()
                        .withBootstrapServers("localhost:9092")
                        .withTopic("file-topic")
                        .withKeyDeserializer(org.apache.kafka.common.serialization.LongDeserializer.class)
                        .withValueDeserializer(org.apache.kafka.common.serialization.StringDeserializer.class))
                .apply(MapElements.into(TypeDescriptor.of(String.class)).via(record -> {
                    String folder = record.getKV().getValue();
                    String path = "file:///tmp/files/" + folder + "/*";
                    System.out.println("path -> " + path);
                    return path;
                }));
        records = absolutePaths
                .apply("read file names", FileIO.matchAll())
                .apply("match file names", FileIO.readMatches())
                .apply("read data from files", TextIO.readFiles());
    } else {
        records = pipeline
                .apply("read from source topic 2", KafkaIO.<Long, String>read()
                        .withBootstrapServers("localhost:9092")
                        .withTopic("data-topic")
                        .withStartReadTime(Instant.parse("2021-09-01T09:30:00-00:00"))
                        .withKeyDeserializer(org.apache.kafka.common.serialization.LongDeserializer.class)
                        .withValueDeserializer(org.apache.kafka.common.serialization.StringDeserializer.class)
                )
                .apply(MapElements.into(TypeDescriptor.of(String.class)).via(record -> record.getKV().getValue()));
    }

    PCollection<String> output = records.apply("counter", new ClientTransformation());
    return output.apply("writing to output topic", KafkaIO.<Void, String>write()
            .withBootstrapServers("localhost:9092")
            .withTopic("output-topic").withValueSerializer(org.apache.kafka.common.serialization.StringSerializer.class)
            .values());

} catch (Exception e) {
    LOG.error("Failed to initialize pipeline due to missing coders", e.getMessage());
    return null;
}

topology in step 1 topology in step 3

0 Answers
Related