You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Flink处理单个Kafka Topic中的键控与批量事件

Flink作业:处理Kafka Topic中的混合事件

我正在开发一个Flink作业,用于处理单个Kafka Topic中的事件流。该事件流包含两类事件:

  • 键控事件:代表Transaction,携带transactionId键,对应插入/更新/删除逻辑;
  • 批量事件:无事务键,需基于业务条件关联所有已处理完成的事务最新状态执行操作。

事件按正确顺序从Topic中到达。

事件类型与处理规则

  • upsert:根据transactionId执行upsert操作
  • delete:根据transactionId执行删除操作
  • paid:将指定agent列表对应的所有交易的支付状态设为true
  • reverse:将指定agent列表对应的所有交易的支付状态设为false
  • publish:将batch_id匹配的所有交易的is_published字段设为true

示例事件明细

event_idevent_typeevent_tstransaction_idbatch_idis_publishedpayments [(agent, paymentstate)]agent
1upsert2024-01-02 16:14:024361false[(1, false),(2, true)]null
2upsert2024-01-02 16:14:044371false[(3, false),(4, false)]null
3upsert2024-01-02 16:14:054361false[(1, true),(2, false)]null
4paid2024-01-02 16:14:10nullnullnullnull[3]
5reverse2024-01-02 16:14:15nullnullnullnull[3]
6publish2024-01-02 16:14:20null1nullnullnull

内容的提问来源于stack exchange,提问作者Valery Maksimenko

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 19:50:24