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;
}