我正在构建一些从 Kafka 读取并使用 Google Cloud Dataflow 写入各种接收器的数据流管道。管道看起来像这样(简化)。
// Example pipeline that writes to BigQuery.
Pipeline.create(options)
.apply(KafkaIO.read().withTopic(options.topic))
.apply(/* Convert to a Row type */)
.setRowSchema(schemaRegistry.lookup(options.topic))
.apply(
BigQueryIO.write<Row>()
.useBeamSchema()
.withCreateDisposition(CreateDispotion.CREATE_IF_NEEDED)
.withProject(options.outputProject)
.withDataset(options.outputDataset)
.withTable(options.outputTable)
)
Run Code Online (Sandbox Code Playgroud)
我计划为我们的每个 Kafka 主题运行一个管道,其中有数百个。管道在规划阶段查找给定主题的模式。这允许BigQueryIO在启动管道之前创建必要的表。
问题:如何在我的 Dataflow 管道中支持不断发展的架构?
我探索了更新现有 Dataflow 作业的选项(使用--update标志)。我的想法是,每当架构发生变化时,我都可以自动执行提交更新作业的过程。但是更新作业似乎会导致大约 3 分钟的停机时间。对于某些工作,那么长时间的停机是行不通的。我正在寻找希望停机时间不超过几秒钟的其他解决方案。
kotlin google-cloud-platform google-cloud-dataflow apache-beam