小编Jam*_*all的帖子

如何编写支持架构演化的流式数据流管道?

我正在构建一些从 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

5
推荐指数
0
解决办法
340
查看次数