ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Airflow 3.1 实时数据管道:用 Flink 和 Kafka 把数据延迟压到分钟级

Airflow 3.1 实时数据管道:用 Flink 和 Kafka 把数据延迟压到分钟级 Airflow 3.1 实时数据管道用 Flink 和 Kafka 把数据延迟压到分钟级【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow数据看板总是比业务动作慢 10 到 20 分钟问题往往不在 Flink 算得慢而在于整条管道被拆成了一串等上一棒跑完再开始的串行任务。Apache Airflow 作为工作流调度平台从 3.1 起将 API 服务器、DAG 处理器、触发器等组件拆分部署配合 Flink 与 Kafka可以搭出一条延迟可控的实时数据管道。这篇不讲概念堆砌只讲怎么把这三者接起来、哪里容易踩坑。先分清两种延迟编排延迟和计算延迟排查实时管道延迟时第一步是把延迟拆成两段看计算延迟Flink 作业内部的处理耗时由并行度、反压、状态后端决定编排延迟任务之间排队、轮询、失败重试带来的等待由调度系统的轮询间隔和任务链长度决定。Airflow 管的是第二段。很多人把 Flink 批作业逻辑塞进一个长任务里假装实时结果 Airflow 的每次 poke、每次重试都直接加在总延迟上。3.1 的组件化架构见下图让 API 服务器、DAG 处理器与触发器可以独立伸缩编排层自身的开销被摊薄这也是它适合承担流式管道指挥角色的原因。正确的分工是Flink 负责长驻的流式计算Airflow 负责部署作业、观察作业状态、处理失败后的修复动作。Airflow 任务不搬运数据只发号施令。用 FlinkKubernetesOperator 部署 Flink 作业如果你用 Kubernetes 跑 Flink通常配合 Flink Kubernetes Operator仓库里 apache.flink provider 提供了一对现成的组件FlinkKubernetesOperator把flinkDeployment自定义资源提交到集群application_file支持本地 yaml/json 文件、YAML 字符串或 JSON 字符串并且支持模板渲染FlinkKubernetesSensor轮询该 Deployment 的状态READY判定成功MISSING、ERROR判定失败还可以把 driver 日志附加到传感器日志里便于排障。典型写法如下作业定义放在 yaml 里DAG 只负责提交 等就绪from airflow.providers.apache.flink.operators.flink_kubernetes import FlinkKubernetesOperator from airflow.providers.apache.flink.sensors.flink_kubernetes import FlinkKubernetesSensor deploy_flink_job FlinkKubernetesOperator( task_iddeploy_flink_job, application_filetemplates/flink_job.yaml, namespacestreaming, ) wait_flink_ready FlinkKubernetesSensor( task_idwait_flink_ready, application_nameorder-aggregation, namespacestreaming, poke_interval15, ) deploy_flink_job wait_flink_ready两个参数值得留意传感器的poke_interval直接决定作业就绪到 Airflow 放行这段延迟的上限15 秒左右是合理区间设成 1 秒只会增加无谓的 API 压力kubernetes_conn_id默认是kubernetes_default你的 Airflow 连接里要配好对应集群的 kube 配置否则任务会卡在鉴权阶段而不是报出明确的错误。Kafka 在这条管道里扮演什么角色Kafka 通常出现在管道两端上游业务事件进 topicFlink 消费后把结果写回 topic 或直接落库。apache.kafka provider 的源码目录里可以看到 hooks、operators、sensors、triggers 等组件对应两类常见需求确认数据到位在 Flink 作业部署之前用 sensor 确认 Kafka 里已经有新数据可消费避免作业空跑异步等待通过 trigger 机制异步等待 Kafka 消息不占用 Airflow worker 进程比长轮询更省资源。如果你的管道是Kafka → Flink → 数据仓库建议把 Kafka topic 的存在性检查放在 DAG 最前端把 Flink 作业的就绪检查放在中段仓库末端再放一个写结果验证任务。这样任一段延迟都能在 UI 里直接定位到对应任务而不是笼统地看到整条管道慢了。配置要点与常见误区安装与依赖两个 provider 独立安装。Flink 这个要求 Airflow ≥ 2.11且依赖apache-airflow-providers-cncf-kubernetes≥ 5.1.0见 provider 文档pip install apache-airflow-providers-apache-flink apache-airflow-providers-apache-kafka几个高频误区指望 Airflow 保证 Flink 内部延迟。Airflow 只能保证作业被正确部署并被监控作业内部反压导致的延迟要去看 Flink 自己的 metrics。批流混用一个 DAG。每天一次的离线重跑和 7×24 的流式守护是两个生命周期硬塞进一个 DAG 会让依赖关系变复杂失败重试互相牵连。传感器轮询间隔调得过低。poke_interval不是越小越好它只影响就绪确认的粒度。忽略application_file的模板化。这个字段支持渲染把日期、topic 名等写死在 yaml 里日后每次变更都要改文件维护成本会快速上升。上线后如何验证延迟管道跑起来后用 Airflow UI 的 landing times 视图对照任务的实际开始时间与计划时间能快速看出是调度排队拖慢的还是作业本身拖慢的排查顺序建议固定为三步看目标任务在 DAG 视图里的开始时间判断是否被上游阻塞看FlinkKubernetesSensor的成功时间与部署任务结束时间的差值即就绪等待耗时都正常的话去 Flink 端查作业的第一条输出时间戳与 Kafka 消息时间戳的差值这才是端到端延迟的主体。下一步想看 Flink 作业经 GCP DataProc 部署的完整示例 DAG可以读 example_dataproc_flink.pyFlink provider 的全部参数说明在 providers/apache/flink/docs/operators.rstKafka 相关组件的介绍从 providers/apache/kafka/README.rst 入手。建议先在测试集群跑通部署 传感器等待这条最小链路再逐步加入 Kafka 前置检查与失败修复任务而不是一开始就搭建完整管道。本文基于项目官方文档与源码编写。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表