【问题标题】:How to trigger airflow jobs based on flink streaming completion for partitions?如何根据分区的 flink 流完成触发气流作业?
【发布时间】:2019-06-06 05:16:15
【问题描述】:

我有一个 flink 流作业,它从 Kafka 读取并写入文件系统中的适当分区。例如,作业被配置为使用写入 /data/date=${date}/hour=${hour} 的存储桶。

如何检测分区已准备好使用,以便相应的气流管道可以在那一小时之上进行一些批处理?

【问题讨论】:

标签: apache-flink airflow flink-streaming lambda-architecture


【解决方案1】:

您可以查看ContinuousFileMonitoringSource 的实现,以了解它如何监控文件系统。然后做一些类似于大卫安德森在你的另一个问题中建议的事情,重新创建一个自定义 ProcessFunction。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-06-28
    • 1970-01-01
    • 2021-05-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多