讲师中心 微信公众号
AI工具推荐 视频效率加速

如何在 Kafka Streams 中跨多个流去重并保留同一流内的重复项

雨敏小哥_4064

雨敏小哥_4064

发布时间:2026-07-26 18:27:23

|

778人浏览过

|

来源于php中文网

原创

如何在 Kafka Streams 中跨多个流去重并保留同一流内的重复项

本文介绍一种基于 kafka streams 的混合处理方案:通过 merge 操作合并多路流,再结合自定义 processor 实现“跨流去重但保留单流内重复”的复杂业务逻辑,解决 outerjoin + aggregate 无法满足的多副本优先级保留需求。

本文介绍一种基于 kafka streams 的混合处理方案:通过 merge 操作合并多路流,再结合自定义 processor 实现“跨流去重但保留单流内重复”的复杂业务逻辑,解决 outerjoin + aggregate 无法满足的多副本优先级保留需求。

在 Kafka Streams 中,标准的 outerJoin 和 aggregate 操作适用于一对一或一对多的关联与聚合场景,但当业务要求区分来源流、保留第二流中所有同值不同键的记录(即“同一原始值在 aug stream 中多次出现,需全部保留”),同时丢弃第一流中对应的所有副本时,传统窗口连接+集合去重的方式会失效——因为 aggregate(LinkedHashSet::new, ...) 会无差别地将所有记录视为等价元素,导致 a:1aug1 和 b:1aug2 被当作重复而仅保留其一。

正确的解法是放弃纯声明式流操作,转而采用 merge() + 自定义 Processor 的组合模式,以获得对每条记录来源、内容、键和顺序的完全控制权。

Upload audio to AIOZ Stream
Upload audio to AIOZ Stream

快速上传音频至 AIOZ Stream API。支持默认或自定义编码配置创建音频对象,上传文件并完成处理后返回音频链接。

下载

✅ 核心思路:Merge 后状态化逐条决策

  1. 统一键映射后 merge:将两路流(raw 和 augmented)各自映射为 (commonKey, CustomRecord) 形式,并调用 mergedStream = rawMapped.merge(augMapped);
  2. 接入自定义 Processor:使用 process() 添加一个有状态的 Processor,内部维护两个关键结构:
    • Map<String, List<CustomRecord>> pendingByCommonKey:暂存按 commonKey 分组的记录(含来源标记);
    • Set<String> seenInAugStream:记录所有已在 augmented stream 中出现过的原始值(如 "1"),用于后续过滤 raw stream 的冲突项;
  3. 按 commonKey 批量提交:当某 commonKey 的所有记录(来自任意流)都到达(可通过时间戳或水印触发 flush),执行业务规则:
    • 若该 key 仅在 raw stream 出现 → 输出 raw 记录;
    • 若该 key 仅在 aug stream 出现 → 输出全部 aug 记录;
    • 若该 key 在两流均出现 → 仅输出 aug stream 的全部记录,彻底忽略 raw stream 对应记录(满足条件 4);
  4. 保序关键:由于原始 raw stream 的顺序需作为最终输出顺序依据,可在 CustomRecord 中额外携带原始 offset 或序列号,并在 Processor 内部对输出结果按此字段排序(或借助 KTable + suppress() 实现事件时间有序输出)。

? 示例 Processor 片段(简化版)

public class DeduplicateAcrossStreamsProcessor 
    implements Processor<String, CustomRecord, String, CustomRecord> {

    private ProcessorContext<String, CustomRecord> context;
    private KeyValueStore<String, List<CustomRecord>> store;

    @Override
    public void init(ProcessorContext<String, CustomRecord> context) {
        this.context = context;
        this.store = context.getStateStore("dedup-store");
    }

    @Override
    public void process(String commonKey, CustomRecord record) {
        // 按 commonKey 聚合记录
        List<CustomRecord> list = store.get(commonKey);
        if (list == null) list = new ArrayList<>();
        list.add(record);
        store.put(commonKey, list);

        // 可选:定时/基于 watermark 触发 flush(此处省略)
        // 或使用 punctuate() 周期性检查并输出
    }

    @Override
    public void punctuate(long timestamp, Punctuator punctuator) {
        // 遍历 store,对每个 commonKey 应用业务规则
        try (KeyValueIterator<String, List<CustomRecord>> iter = store.all()) {
            while (iter.hasNext()) {
                KeyValue<String, List<CustomRecord>> entry = iter.next();
                List<CustomRecord> records = entry.value;

                List<CustomRecord> augOnly = records.stream()
                    .filter(r -> r.origin == OriginStream.AUGMENTED)
                    .collect(Collectors.toList());

                if (!augOnly.isEmpty()) {
                    // 条件 2 & 4:只要存在 aug 记录,就全部输出,且不输出任何 raw
                    augOnly.forEach(r -> context.forward(r.key(), r));
                } else {
                    // 条件 1 & 3:仅 raw 或仅 aug(但此处 augOnly 为空,故只剩 raw)
                    records.stream()
                        .filter(r -> r.origin == OriginStream.RAW)
                        .forEach(r -> context.forward(r.key(), r));
                }
                store.delete(entry.key); // 清理已处理 key
            }
        }
    }
}

⚠️ 注意事项与最佳实践

  • 状态存储必须启用:store 需在 Topology 中显式定义为 Stores.persistentKeyValueStore(...),并绑定到 Processor;
  • 容错与恢复:确保 CustomRecord 可序列化,且 Processor 的 punctuate() 逻辑幂等(如使用 store.delete() 配合 context.commit());
  • 顺序保证:若 raw stream 的物理顺序至关重要,建议在 CustomRecord 中嵌入 rawOffset 字段,并在 punctuate() 输出前按该字段排序;
  • 性能考量:避免在 process() 中做耗时操作;高频 key 可考虑分片或 TTL 策略防止状态无限增长;
  • 替代方案提示:对于低吞吐、强一致性场景,也可先将两流分别写入 KTable,再用 transform() 查询 aug 表是否存在对应 key,但该方式延迟更高且无法天然保留 aug 流内多键副本。

综上,当 Kafka Streams 的 DSL 层无法表达“按来源流差异化去重”这类复杂语义时,转向 Processor API 并辅以状态存储,是最直接、可控且可验证的工程解法。

相关文章

Kafka Eagle可视化工具
Kafka Eagle可视化工具

Kafka Eagle是一款结合了目前大数据Kafka监控工具的特点,重新研发的一块开源免费的Kafka集群优秀的监控工具。它可以非常方便的监控生产环境中的offset、lag变化、partition分布、owner等,有需要的小伙伴快来保存下载体验吧!

下载

相关标签:

本站声明:本文内容由网友自发贡献,版权归原作者所有,本站不承担相应法律责任。如您发现有涉嫌抄袭侵权的内容,请联系admin@php.cn

热门AI工具

更多
墨刀AI
墨刀AI Hot

一款AI图像与设计工具,主要用于产品经理的专属智能体,适合需要提升相关任务效率的用户。

豆包大模型

豆包大模型是一款由字节跳动推出的企业级大语言模型服务平台。

PixPix
PixPix Hot

PixPix是一款面向电商视觉生产的AI商品图生成工具。

PixTV
PixTV Hot

PixTV是一款面向AIGC内容创作的AI视频生成工具。

DeepSeek

DeepSeek是一款面向对话、写作、编程和推理场景的AI大模型工具。

WorkBuddy

一款AI办公效率工具,主要用于腾讯云推出的AI原生桌面智能体工作台,适合需要提升相关任务效率的用户。

超级简历WonderCV

一款AI办公效率工具,主要用于免费求职简历模版下载制作,应届生职场人必备简历制作神器,适合需要提升相关任务效率的用户。

切问学术

切问学术是一款AI论文写作工具,复旦大学NLP团队推出的AI学术智能体。

讯飞智作

讯飞智作是一款AI视频创作工具,AI文本配音工具,数字人课程、营销视频制作。

相关专题

更多
kafka消费者组有什么作用
kafka消费者组有什么作用

kafka消费者组的作用:1、负载均衡;2、容错性;3、广播模式;4、灵活性;5、自动故障转移和领导者选举;6、动态扩展性;7、顺序保证;8、数据压缩;9、事务性支持。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

2366

2024.01.12

kafka消费组的作用是什么
kafka消费组的作用是什么

kafka消费组的作用:1、负载均衡;2、容错性;3、灵活性;4、高可用性;5、扩展性;6、顺序保证;7、数据压缩;8、事务性支持。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

570

2024.02.23

rabbitmq和kafka有什么区别
rabbitmq和kafka有什么区别

rabbitmq和kafka的区别:1、语言与平台;2、消息传递模型;3、可靠性;4、性能与吞吐量;5、集群与负载均衡;6、消费模型;7、用途与场景;8、社区与生态系统;9、监控与管理;10、其他特性。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

544

2024.02.23

Java 流式处理与 Apache Kafka 实战
Java 流式处理与 Apache Kafka 实战

本专题专注讲解 Java 在流式数据处理与消息队列系统中的应用,系统讲解 Apache Kafka 的基础概念、生产者与消费者模型、Kafka Streams 与 KSQL 流式处理框架、实时数据分析与监控,结合实际业务场景,帮助开发者构建 高吞吐量、低延迟的实时数据流管道,实现高效的数据流转与处理。

590

2026.02.04

LLVM自定义Pass怎么写
LLVM自定义Pass怎么写

本专题聚焦LLVM自定义Pass开发,整理Pass类结构、run()方法、PreservedAnalyses、CMake构建、插件注册、-load-pass-plugin加载和测试用例编写流程。

20

2026.09.30

LLVM RISC-V参数配置教程
LLVM RISC-V参数配置教程

本专题介绍LLVM对RISC-V基础ISA和扩展的支持方式,涵盖RV32、RV64、标准扩展、实验性扩展、厂商扩展、-menable-experimental-extensions和版本差异。

0

2026.09.30

LLVM IR中间表示入门指南
LLVM IR中间表示入门指南

本专题整理LLVM IR的核心概念,包括中间表示作用、模块结构、函数、基本块、SSA形式、类型系统和常见语法,帮助新手理解LLVM编译流程中的关键层。

0

2026.09.30

PDF转图片方法
PDF转图片方法

需要把 PDF 页面用于上传、预览、分享或图片归档时,PDF 转图片方法专题整理 JPG/PNG 格式选择、逐页导出、清晰度设置、批量下载和结果检查等流程,帮助用户稳定完成 PDF 图片化处理。

20

2026.09.30

PixTV AI视频生成与无限画布创作
PixTV AI视频生成与无限画布创作

PixTV专题整理AI视频与视觉内容创作相关功能使用教程,涵盖AI生图、视频生成、无限画布、多模型创作、素材管理、声音音乐及视频剪辑等功能,帮助用户快速掌握PixTV从创意到成片的完整制作方法。

20

2026.09.29

热门下载

更多
网站特效
/
网站源码
/
网站素材
/
前端模板

精品课程

更多
相关推荐
/
热门推荐
/
最新课程
CentOS 官方文档
CentOS 官方文档

共0课时 | 0人学习

极客学院Java8新特性视频教程
极客学院Java8新特性视频教程

共17课时 | 4.3万人学习

关于我们 免责申明 举报中心 意见反馈 讲师合作 广告合作 最新更新
php中文网:公益在线php培训,帮助PHP学习者快速成长!
关注服务号
PHP中文网订阅号
每天精选资源文章推送

Copyright 2014-2026 https://www.php.cn/ All Rights Reserved | php.cn