如何更新运行中的Apache Beam/Dataflow作业 动态调整Kafka主题不丢数据
GCP Dataflow多Kafka主题动态更新作业无数据丢失方案咨询
场景说明
在Google Cloud Platform环境中,我们使用Dataflow将Kafka中的消息接入Google Cloud Storage,上传的对象按小时分区,存储路径格式如下:
gs://output-bucket/topics/<topic>/year=YYYY/month=MM/day=dd/hour=HH/<windowStart>-<windowEnd>-<pane>.<format>
晚到数据会落入额外的pane中存储。
初期方案为单个Dataflow作业通过KafkaIO接入多个Kafka主题,该方案运行状态正常。但业务侧需要不定期增删接入列表中的Kafka主题并同步更新作业配置,更新过程需保证无数据丢失(包含尚未写入Cloud Storage的窗口内数据),同时保留晚到数据的处理状态。
预期更新效果
- 调整接入主题列表,完成新增/删除主题操作
- 使用新的主题参数更新运行中的Dataflow作业,无需重启全量任务
- 原有主题继续从当前消费offset接入,新增主题从最早offset开始接入
现存问题
参考官方管道更新指南操作,仅修改输入主题参数更新现有管道时,出现如下报错:
CESTError message from worker: java.lang.IllegalStateException: checkPointMark and assignedPartitions should match
经测试验证,原生更新方式不支持该类配置变更,核心诉求为:单个Dataflow作业读取多个Kafka主题时,如何正确更新现有管道的主题参数,满足上述无数据丢失的动态调整需求?
现有基础管道代码
pipeline .apply( "ReadFromKafka", KafkaIO.<String, String>read() .withBootstrapServers(options.getBootstrapServers()) .withTopics(parseTopics(options.getInputTopics())) .withCreateTime(Duration.standardMinutes(10)) .withReadCommitted() .commitOffsetsInFinalize() .withKeyDeserializer(StringDeserializer.class) .withValueDeserializer(StringDeserializer.class) .withConsumerConfigUpdates(parseConsumerConfig(options))) .apply("CreateWindows", HourlyWindowsWithAllowedLateness.of(Duration.standardDays(7))) .apply( "WriteFileOutput", FileIO.<String, KafkaRecord<String, String>>writeDynamic() .by(DestinationFunction.of()) .via( Contextful.fn(KafkaRecordToGenericRecordFunction.of()), AvroIO.sink(KafkaRecordToGenericRecordFunction.SCHEMA)) .withDestinationCoder(StringUtf8Coder.of()) .withNumShards(1) .withTempDirectory(options.getTempDirectory()) .to(options.getOutputDirectory()) .withNaming(NamingFunction.avro()));
内容的提问来源于stack exchange,提问作者Marcel Neumann
相关产品推荐
相关产品推荐

