如何用Flink处理单个Kafka Topic中的键控与批量事件
Flink作业:处理Kafka Topic中的混合事件
我正在开发一个Flink作业,用于处理单个Kafka Topic中的事件流。该事件流包含两类事件:
- 键控事件:代表Transaction,携带transactionId键,对应插入/更新/删除逻辑;
- 批量事件:无事务键,需基于业务条件关联所有已处理完成的事务最新状态执行操作。
事件按正确顺序从Topic中到达。
事件类型与处理规则
upsert:根据transactionId执行upsert操作delete:根据transactionId执行删除操作paid:将指定agent列表对应的所有交易的支付状态设为truereverse:将指定agent列表对应的所有交易的支付状态设为falsepublish:将batch_id匹配的所有交易的is_published字段设为true
示例事件明细
| event_id | event_type | event_ts | transaction_id | batch_id | is_published | payments [(agent, paymentstate)] | agent |
|---|---|---|---|---|---|---|---|
| 1 | upsert | 2024-01-02 16:14:02 | 436 | 1 | false | [(1, false),(2, true)] | null |
| 2 | upsert | 2024-01-02 16:14:04 | 437 | 1 | false | [(3, false),(4, false)] | null |
| 3 | upsert | 2024-01-02 16:14:05 | 436 | 1 | false | [(1, true),(2, false)] | null |
| 4 | paid | 2024-01-02 16:14:10 | null | null | null | null | [3] |
| 5 | reverse | 2024-01-02 16:14:15 | null | null | null | null | [3] |
| 6 | publish | 2024-01-02 16:14:20 | null | 1 | null | null | null |
内容的提问来源于stack exchange,提问作者Valery Maksimenko
相关产品推荐
相关产品推荐

