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

如何在 Kafka Streams 中实现跨流去重并保留同流内重复项

夏萱小哥_4054

夏萱小哥_4054

发布时间:2026-07-26 12:41:07

|

484人浏览过

|

来源于php中文网

原创

如何在 Kafka Streams 中实现跨流去重并保留同流内重复项

本文介绍一种基于 kafka streams 的混合处理方案:通过 merge 操作合并多路数据流,再结合自定义 processor 实现“跨流去重但保留单流内重复”的业务逻辑,精准满足复杂消息路由与去重需求。

本文介绍一种基于 kafka streams 的混合处理方案:通过 merge 操作合并多路数据流,再结合自定义 processor 实现“跨流去重但保留单流内重复”的业务逻辑,精准满足复杂消息路由与去重需求。

在 Kafka Streams 应用中,常见的 reduce 或 outerJoin 无法直接支持“同一键下保留第二流所有副本、同时丢弃第一流对应记录”这一非幂等性去重逻辑——因为 reduce 是两两归约,join 是成对匹配,二者均无法感知某键在第二流中出现的次数及全部实例。

正确解法是放弃纯 DSL(Domain Specific Language)方式,转而采用 KStream#merge() + 自定义 Processor 的混合编程模型:

  1. 先合并流:将原始流(无键、带时间戳)与增强流(有键、值变形)统一映射为相同键结构后合并;
  2. 再状态化处理:在自定义 Processor 中维护两个状态:
    • Set<String> 记录所有在第二流(augmented)中出现过的去重键(用于识别“该键是否来自第二流”);
    • List<CustomMessageDetailsWithKeyAndOrigin> 缓存当前键的所有第二流消息(支持保留多份);
  3. 按需输出:遍历合并后的每条记录,依据其来源(RAW / AUGMENTED)和全局键状态,执行差异化路由:
    • 若为 AUGMENTED:加入缓存列表,并标记该键“已见于第二流”;
    • 若为 RAW:仅当该键未在第二流中出现过时才转发;
  4. 最终排序保障:因原始流顺序需保留,建议在 merge 前为原始流添加单调递增序列号(如 mapValues((v, idx) -> new EnrichedValue(v, idx))),并在 Processor 中按此序号做缓冲排序(或依赖 Kafka 分区局部有序 + 后续 Flink/Spark 二次排序)。

以下是关键代码片段示例:

Upload audio to AIOZ Stream
Upload audio to AIOZ Stream

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

下载
// 步骤1:统一键映射并合并
KStream<String, CustomMsg> mappedRaw = rawInputStream
    .map((k, v) -> KeyValue.pair(getCommonKey(v), 
        new CustomMsg(v, "", OriginStream.RAW)));

KStream<String, CustomMsg> mappedAug = augmentedInputStream
    .map((k, v) -> KeyValue.pair(getCommonKey(v), 
        new CustomMsg(v, k, OriginStream.AUGMENTED)));

KStream<String, CustomMsg> merged = mappedRaw.merge(mappedAug);

// 步骤2:接入自定义 Processor(需注册 StateStore)
merged.process(() -> new DedupProcessor(), "dedup-store");

自定义 DedupProcessor 核心逻辑(简化版):

public class DedupProcessor implements Processor<String, CustomMsg, String, CustomMsg> {
    private ProcessorContext<String, CustomMsg> context;
    private KeyValueStore<String, List<CustomMsg>> augStore; // 存储第二流全量消息
    private KeyValueStore<String, Boolean> seenInAugStore;    // 标记键是否出现在第二流

    @Override
    public void init(ProcessorContext<String, CustomMsg> context) {
        this.context = context;
        this.augStore = context.getStateStore("aug-store");
        this.seenInAugStore = context.getStateStore("seen-flag-store");
    }

    @Override
    public void process(String key, CustomMsg value) {
        if (value.origin == OriginStream.AUGMENTED) {
            // 保存第二流消息(允许重复 key)
            List<CustomMsg> list = augStore.get(key);
            if (list == null) list = new ArrayList<>();
            list.add(value);
            augStore.put(key, list);
            seenInAugStore.put(key, true); // 标记该键存在第二流
        } else { // RAW 流
            if (!Boolean.TRUE.equals(seenInAugStore.get(key))) {
                // 仅当该 key 从未出现在第二流时才转发原始消息
                context.forward(key, value, To.all());
            }
        }
    }

    @Override
    public void punctuate(long timestamp) {
        // 定期遍历 augStore,将每个 key 对应的全部第二流消息逐条 forward
        try (KeyValueIterator<String, List<CustomMsg>> iter = augStore.all()) {
            while (iter.hasNext()) {
                KeyValue<String, List<CustomMsg>> entry = iter.next();
                for (CustomMsg msg : entry.value) {
                    context.forward(entry.key, msg, To.all());
                }
            }
            augStore.flush(); // 清空已处理批次(按需设计清理策略)
        }
    }
}

⚠️ 注意事项:

  • 状态存储需配置 TTL(如 TimeWindowedStore)防止无限增长;
  • punctuate() 触发时机影响实时性,建议结合 WallclockTimeExtractor 或事件时间窗口控制;
  • 若要求严格保序,应在 CustomMsg 中嵌入原始流的序列号或时间戳,并在 punctuate 阶段排序后输出;
  • 生产环境务必启用 enable.auto.commit 和 processing.guarantee=exactly_once_v2 保障端到端精确一次语义。

该方案突破了 Kafka Streams DSL 的表达边界,以可控的状态管理换取业务逻辑的完全自主权,是处理“条件性多副本保留”类场景的稳健实践路径。

相关文章

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

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

下载

相关标签:

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

热门AI工具

更多
切问学术

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

WorkBuddy

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

墨刀AI
墨刀AI Hot

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

豆包大模型

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

Lovart
Lovart Hot

一款面向视觉设计创作的AI设计平台,可通过智能体和画布工作流辅助制作海报、Logo、网页、PPT及其他视觉内容。

超级简历WonderCV

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

咔片AIPPT

一款在线AI演示文稿制作工具,可根据主题和内容需求辅助生成PPT结构与页面,提高演示材料制作效率。

讯飞智作

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

DeepSeek

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

相关专题

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

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

2326

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

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

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

0

2026.09.30

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

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

0

2026.09.29

Buffalo框架数据库开发全教程
Buffalo框架数据库开发全教程

本专题围绕Buffalo框架数据库开发,讲解database.yml多环境配置、soda与fizz迁移生成回滚、模型结构体标签、增删改查与条件查询、一对多与多对多关联、数据校验、回调钩子、事务处理及原生SQL执行能力。

200

2026.09.23

Buffalo框架路由与请求处理实操指南
Buffalo框架路由与请求处理实操指南

本专题讲解Buffalo框架路由与请求处理机制,涵盖路由注册与分组、资源路由、Handler编写规范、Context上下文方法、参数绑定、中间件编写挂载、Session与Cookie读写、Flash消息及错误页面定制方法。

120

2026.09.23

Buffalo框架零基础入门教程
Buffalo框架零基础入门教程

本专题整理Buffalo框架入门内容,涵盖Go环境准备、buffalo CLI安装、新项目生成、目录结构说明、dev热加载启动、数据库连接配置与常见报错排查,帮助新手按约定优于配置的思路跑通第一个Buffalo框架应用。

100

2026.09.23

热门下载

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

精品课程

更多
相关推荐
/
热门推荐
/
最新课程
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