在 Apache Flink 中消费带键 Kafka 记录的实践教程

在 Apache Flink 中消费带键 Kafka 记录的实践教程

本教程旨在指导您如何在 apache flink 中高效消费带有键的 kafka 记录。文章详细介绍了使用自定义 `kafkarecorddeserializationschema` 来解析 kafka `consumerrecord` 中的键、值、时间戳等信息,并提供了完整的 flink 应用程序代码示例。通过遵循本文的步骤,您可以轻松地构建能够处理复杂 kafka 消息结构的 flink 流处理应用。

1. 理解带键 Kafka 记录及其重要性

在 Kafka 中,消息(记录)通常包含一个可选的键(Key)和一个值(Value)。键在许多场景下都至关重要,例如:

消息顺序保证:同一个键的所有消息会被发送到同一个分区,从而保证了这些消息的消费顺序。状态管理:在 Flink 等流处理框架中,键是进行有状态操作(如聚合、连接)的基础。数据路由:消费者可以根据键来过滤或路由消息。

当使用 kafka-console-producer.sh 并指定 –property “parse.key=true” –property “key.separator=:” 时,生产者会从输入中解析出键和值,并将它们作为独立的字段发送到 Kafka。例如,myKey:myValue 会被解析为键 myKey 和值 myValue。

2. Flink KafkaSource 的默认行为与限制

Apache Flink 提供了 KafkaSource 作为消费 Kafka 数据的首选连接器。然而,当您使用 KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class) 这样的默认配置时,KafkaSource 仅会反序列化 Kafka 记录的值部分,而忽略其键、时间戳、分区、偏移量以及头部信息。这对于只需要处理消息值的场景是足够的,但对于需要访问键或其它元数据的应用来说,这种方式就显得力不从心。

以下是仅读取非带键记录的示例代码,它无法获取 Kafka 记录的键:

import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.connector.kafka.source.KafkaSource;import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;import org.apache.kafka.common.serialization.StringDeserializer;public class FlinkValueOnlyKafkaConsumer {    public static void main(String[] args) throws Exception {        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();        String bootstrapServers = "localhost:9092"; // 替换为您的Kafka地址        KafkaSource source = KafkaSource.builder()                .setBootstrapServers(bootstrapServers)                .setTopics("test3")                .setGroupId("1")                .setStartingOffsets(OffsetsInitializer.earliest())                .setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class))                .build();        DataStream stream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source");        stream.map((MapFunction) value -> "Receiving from Kafka : " + value).print();        env.execute("Flink Value-Only Kafka Consumer");    }}

3. 自定义 KafkaRecordDeserializationSchema 读取带键记录

要从 Kafka 记录中获取键、值、时间戳等所有信息,您需要实现一个自定义的 KafkaRecordDeserializationSchema。这个接口的 deserialize 方法会接收一个 ConsumerRecord 对象,该对象提供了对原始字节形式的键、值、时间戳、分区、偏移量以及头部信息的完全访问。

3.1 定义自定义反序列化器

首先,创建一个实现 KafkaRecordDeserializationSchema 接口的类。在这个示例中,我们将反序列化键和值都为 String 类型,并将它们与时间戳一起封装到一个 Tuple3 对象中输出。

PicDoc PicDoc

AI文本转视觉工具,1秒生成可视化信息图

PicDoc 6214 查看详情 PicDoc

import org.apache.flink.api.common.serialization.DeserializationSchema;import org.apache.flink.api.common.typeinfo.TypeInformation;import org.apache.flink.api.java.tuple.Tuple3;import org.apache.flink.util.Collector;import org.apache.kafka.clients.consumer.ConsumerRecord;import org.apache.kafka.common.serialization.StringDeserializer;import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;import java.io.IOException;/** * 自定义 Kafka 记录反序列化器,用于解析键、值和时间戳。 * 输出类型为 Tuple3 */public class KeyedKafkaRecordDeserializationSchema implements KafkaRecordDeserializationSchema<Tuple3> {    // transient 关键字确保这些反序列化器不会被 Flink 的序列化机制尝试序列化    private transient StringDeserializer keyDeserializer;    private transient StringDeserializer valueDeserializer;    /**     * 在反序列化器初始化时调用,用于设置内部状态。     * 通常在这里初始化 Kafka 客户端的反序列化器。     */    @Override    public void open(DeserializationSchema.InitializationContext context) throws Exception {        // 根据 Kafka 生产者实际使用的序列化器来选择这里的反序列化器        // 假设键和值都是字符串,使用 StringDeserializer        keyDeserializer = new StringDeserializer();        valueDeserializer = new StringDeserializer();    }    /**     * 核心反序列化逻辑。     *     * @param record Kafka 原始的 ConsumerRecord 对象,包含字节数组形式的键和值。     * @param out    用于收集反序列化结果的 Collector。     * @throws IOException 如果反序列化过程中发生 I/O 错误。     */    @Override    public void deserialize(ConsumerRecord record, Collector<Tuple3> out) throws IOException {        // 反序列化键        String key = (record.key() != null) ? keyDeserializer.deserialize(record.topic(), record.key()) : null;        // 反序列化值        String value = (record.value() != null) ? valueDeserializer.deserialize(record.topic(), record.value()) : null;        // 获取时间戳        long timestamp = record.timestamp();        // 将反序列化后的键、值和时间戳封装成 Tuple3 并发出        out.collect(new Tuple3(key, value, timestamp));    }    /**     * 返回此反序列化器生产的数据类型信息。     * Flink 使用此信息进行类型检查和序列化。     */    @Override    public TypeInformation<Tuple3> getProducedType() {        // 使用 TypeHint 来获取泛型类型信息        return TypeInformation.of(new org.apache.flink.api.java.typeutils.TypeHint<Tuple3>() {});    }}

注意事项:

open 方法:在反序列化器首次使用时调用,用于初始化资源。将 Kafka 客户端的反序列化器(如 StringDeserializer)放在这里初始化可以避免在每次 deserialize 调用时重复创建对象,提高效率。deserialize 方法:这是核心逻辑所在。ConsumerRecord 提供了 key()、value()、timestamp()、topic()、partition()、offset() 和 headers() 等方法。您可以使用 Kafka 客户端提供的反序列化器(例如 StringDeserializer、LongDeserializer 或自定义的 Avro/Protobuf 反序列化器)来将 byte[] 转换为实际的数据类型。getProducedType 方法:必须返回此反序列化器将发出的数据流的 TypeInformation。这对于 Flink 的类型系统至关重要。

3.2 在 Flink KafkaSource 中使用自定义反序列化器

接下来,将我们自定义的 KeyedKafkaRecordDeserializationSchema 应用到 KafkaSource 中:

import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.java.tuple.Tuple3;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.connector.kafka.source.KafkaSource;import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;public class FlinkKeyedKafkaConsumer {    public static void main(String[] args) throws Exception {        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();        String bootstrapServers = "localhost:9092"; // 替换为您的 Kafka 地址        String topic = "test3";        String groupId = "1";        // 构建 KafkaSource,并指定我们自定义的反序列化器        KafkaSource<Tuple3> source = KafkaSource.<Tuple3>builder()                .setBootstrapServers(bootstrapServers)                .setTopics(topic)                .setGroupId(groupId)                .setStartingOffsets(OffsetsInitializer.earliest())                .setDeserializer(new KeyedKafkaRecordDeserializationSchema()) // 使用自定义反序列化器                .build();        // 从 KafkaSource 创建数据流        DataStream<Tuple3> stream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Keyed Kafka Source");        // 对数据流进行操作,现在可以访问键、值和时间戳        stream.map(record -> "Key: " + record.f0 + ", Value: " + record.f1 + ", Timestamp: " + record.f2)              .print();        // 执行 Flink 作业        env.execute("Flink Keyed Kafka Consumer");    }}

3.3 Kafka 生产者示例(用于测试)

为了测试上述 Flink 消费者,您可以使用以下命令启动一个 Kafka 控制台生产者,它会生成带键的记录:

bin/kafka-console-producer.sh --topic test3 --property "parse.key=true" --property "key.separator=:" --bootstrap-server localhost:9092

然后,在控制台中输入 myKey:myValue 这样的消息,Flink 消费者将能够正确解析出 myKey 作为键,myValue 作为值。

4. 总结

通过实现自定义的 KafkaRecordDeserializationSchema,您可以完全控制 Flink 如何从 Kafka 的原始 ConsumerRecord 中提取和反序列化数据。这不仅限于键和值,还可以包括时间戳、主题、分区、偏移量甚至自定义头部信息。这种灵活性使得 Flink 能够处理各种复杂的 Kafka 消息格式,为构建强大的流处理应用提供了坚实的基础。在实际应用中,请确保自定义反序列化器中使用的 Kafka 客户端反序列化器与生产者使用的序列化器保持一致。

以上就是在 Apache Flink 中消费带键 Kafka 记录的实践教程的详细内容,更多请关注创想鸟其它相关文章!

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。
如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 chuangxiangniao@163.com 举报,一经查实,本站将立刻删除。
发布者:程序猿,转转请注明出处:https://www.chuangxiangniao.com/p/988581.html

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
sql 中 between 用法_sql 中 between 范围查询技巧
上一篇 2025年12月1日 21:12:05
如何使用CSS浮动实现图文混排效果_实战案例解析
下一篇 2025年12月1日 21:12:06

相关推荐

  • PHP简易路由框架构建:从URL解析到动态控制器加载的实践指南

    本文旨在指导读者构建一个基础的PHP路由系统,实现URL路径到控制器方法的高效映射。内容涵盖URL解析、控制器动态加载、方法调用以及关键的错误处理机制,特别强调如何避免常见的“未定义变量”错误和文件包含路径问题,确保路由系统稳定且易于维护。 一、路由系统核心原理 构建一个简单的php路由系统,其核心…

    2026年9月21日
    100
  • 如何在Java中理解Java I/O与NIO机制

    传统I/O是阻塞式流模型,适用于低并发场景;NIO基于缓冲区与通道,支持非阻塞和多路复用,适合高并发网络应用,核心区别在于线程模型与资源利用率。 Java中的I/O(输入/输出)与NIO(New I/O)是处理数据读写的核心机制,理解它们的区别和使用场景对开发高性能应用至关重要。传统I/O基于流模型…

    2026年9月21日
    000
  • UC浏览器网页上的文字无法选中复制怎么办 UC浏览器解决网页文字禁止复制问题

    答案:可通过开发者工具、阅读模式、打印预览、OCR识别或自定义脚本解除UC浏览器网页复制限制。具体操作依次为:开启开发者工具并执行JavaScript代码解除限制;启用阅读模式净化页面内容;使用打印预览重新渲染页面以选中文字;对截图应用OCR技术提取文本;添加书签脚本自动移除禁用选择的代码,从而实现…

    2026年9月21日
    000
  • JavaScript中的尾调用优化(TCO)在ES6中如何工作?

    尾调用是指函数的最后一个动作调用另一个函数,ES6引入尾调用优化以重用栈帧、避免内存溢出,支持真正的尾递归,如阶乘函数通过累积参数实现。 尾调用优化(Tail Call Optimization, TCO)是ES6引入的一项语言特性,目的是在特定条件下重用函数调用栈帧,避免不必要的内存增长,从而支持…

    2026年9月21日
    100
  • 为什么iPhoneSE2022屏幕无响应如何强制重启?快速按音量键后长按电源键

    首先尝试强制重启,若无效则检查充电状态,最后可通过恢复模式重装系统。具体为:1. 按音量+、音量-后长按电源键10秒以上;2. 充电15分钟观察是否响应;3. 连电脑进入恢复模式恢复系统。 如果您尝试唤醒或操作您的iPhone SE(2022款),但屏幕无响应或显示黑屏,可能是系统临时卡死或软件冲突…

    2026年9月21日
    000
  • 抖音蝴蝶号无人直播带货操作流程及注意事项

    抖音蝴蝶号无人直播带货操作流程及注意事项抖音蝴蝶号无人直播带货操作流程及注意事项抖音蝴蝶号无人直播带货操作流程及注意事项抖音蝴蝶号无人直播带货操作流程及注意事项

    “抖音蝴蝶号无人直播带货”是一种通过自动化或半自动化技术实现的直播销售模式。①其核心在于摆脱真人主播限制,实现24小时不间断直播,提升效率与流量利用率;②关键步骤包括明确账号定位与商品选择、准备高质量且丰富的内容素材、利用虚拟人或预录内容实现直播推流、结合智能客服模拟评论区互动;③优势在于降低人力成…

    2026年9月21日 用户投稿
    500
  • Java语法基础有哪些新手必学的核心知识

    掌握Java基本数据类型与变量声明,如int、double、char和boolean,并理解强类型语言特性;2. 熟悉运算符与表达式,包括算术、比较和逻辑运算符,奠定程序逻辑基础。 Java语法基础是每个初学者必须掌握的内容,只有打好根基,才能顺利进阶面向对象编程和实际项目开发。以下是新手必学的核心…

    2026年9月21日
    200
  • 音乐文件占用空间太多怎么办_音乐文件占用空间太多如何整理详细指南

    解决音乐文件占空间问题的关键是压缩与整理:先用软件或在线工具降低比特率压缩体积,再按场景分类、利用元数据自动归集,并通过听歌片段和BPM判断保留内容,避免重复与误删。 音乐文件占空间太多,核心解决办法就两条:一是压缩单个文件体积,二是通过有效分类管理提升使用效率。直接删歌不是长久之计,学会整理和优化…

    2026年9月21日
    000
  • 升级X86架构性能大提升!极空间Z2 Ultra图赏

    升级X86架构性能大提升!极空间Z2 Ultra图赏升级X86架构性能大提升!极空间Z2 Ultra图赏升级X86架构性能大提升!极空间Z2 Ultra图赏升级X86架构性能大提升!极空间Z2 Ultra图赏

    10月23日,极空间正式推出全新双盘位nas产品——极空间z2 ultra,官方售价为1899元,参与国家补贴后仅需1457元,性价比进一步提升。 此次发布的Z2 Ultra最大的亮点在于采用X86架构处理器,相较以往使用的ARM平台,性能实现飞跃式提升,运行速度显著加快。更重要的是,新架构对Doc…

    2026年9月21日 用户投稿
    200
  • 数据库分库分表(Sharding)策略

    在现代应用程序中,随着数据量的增长,单一数据库的性能和容量往往难以满足需求。这时,数据库分库分表(Sharding)策略就成了一个关键的解决方案。那么,如何设计和实现一个有效的分库分表策略呢?让我们深入探讨一下。 在我的职业生涯中,我曾多次参与大型项目的数据库优化,其中分库分表是常见的挑战之一。我记…

    2026年9月21日
    000
  • 如何在Java中实现个人财务管理工具

    首先设计Transaction、FinanceManager和Budget核心类,实现交易记录、统计分析与预算控制功能,通过ArrayList管理数据,使用LocalDate处理日期,结合ObjectOutputStream持久化存储,初期采用Scanner构建控制台菜单实现增删查改与报表展示,后期…

    2026年9月21日
    000
  • X旗下Grok上线即时语音搜索,挑战Google引领搜索新方向

    近日,x平台旗下的ai助手grok正式推出了“即时语音搜索”功能。用户现在可以通过语音直接提问,触发实时网页检索,并迅速获得整合后的精准答案。此举意在优化信息获取流程,推动人机交互向更自然、高效的方向演进。 该语音搜索模式实现了“即说即搜即答”的流畅体验。例如,当用户提出“星舰发射的具体时间是什么?…

    2026年9月21日
    100
  • Laravel应用的安全审计(Security Audit)方法

    进行安全审计对laravel应用至关重要,因为它能发现并修复安全漏洞,提升整体安全性和用户信任度。具体方法包括:1. 代码审查,确保无未过滤输入和弱密码;2. 配置文件安全性,保护敏感信息;3. 依赖管理,更新第三方包;4. 用户认证和授权,防止未授权访问;5. 日志和监控,检测异常行为。 在讨论L…

    2026年9月21日
    100
  • Laravel 8 登录后重定向到仪表盘的全面指南

    本文深入探讨了 Laravel 8 中用户登录后重定向到仪表盘的多种策略。我们将详细解析默认的重定向机制,包括 LoginController 和 RedirectIfAuthenticated 中间件,并重点介绍如何通过自定义登录逻辑实现精确的重定向控制,同时提供示例代码和常见问题排查建议,确保用…

    2026年9月21日
    000
  • Guava Multimap:高效获取并打印指定键的所有关联值

    guava multimap是处理一键多值映射关系的强大工具。要获取特定键的所有关联值,应直接使用其提供的`multimap#get(k)`方法。该方法会返回一个包含所有匹配值的`collection`,即使键不存在,也会返回一个空集合而非`null`,从而简化了值检索和空值处理逻辑,是比手动迭代键…

    2026年9月21日
    000
  • 控制台命令(Console Command)开发

    控制台命令是程序员日常工作中不可或缺的工具,它提高了开发效率并帮助理解和控制程序运行。1) 通过简单的文本输入,完成复杂任务,如文件管理和系统监控。2) 控制台命令可用于快速调试、测试代码和自动化重复工作。3) 开发控制台命令时需注意安全性和兼容性问题。4) 控制台命令可实现有趣功能,如监控服务器资…

    2026年9月21日
    100
  • 链路追踪(OpenTelemetry/Jaeger)集成

    要将opentelemetry和jaeger集成到java应用中,需按以下步骤操作:1.配置jaeger exporter,2.初始化opentelemetry,3.创建并管理span。通过这种方式,你可以有效地追踪和分析微服务间的调用链路,提升系统性能。 在现代微服务架构中,链路追踪已经成为诊断和…

    2026年9月21日
    000
  • Maingear电脑黑屏问题如何修复?专业级主机BIOS设置方法详尽

    Maingear电脑黑屏问题通常由BIOS设置、硬件接触不良或显示输出配置引起。首先应尝试进入BIOS,检查并调整显卡输出模式为PCIe/PEG,确保未误设为集成显卡;排查PCIe插槽模式兼容性,必要时切换为Gen3或Auto;若启动异常,可尝试切换UEFI/Legacy模式或恢复BIOS默认设置(…

    2026年9月21日
    000
  • 实测!Sora 2长视频优势大,Vidu Q2细节处理更胜一筹

    近日,AI视频工具领域的竞争愈发激烈。OpenAI推出的Sora 2刚刚登顶美区App Store榜单,国产新秀Vidu Q2便携重磅升级版本强势入局,引发广泛关注。不少从事自媒体创作与影视剪辑的朋友都在思考:这两款AI视频生成器,究竟谁更胜一筹?出于好奇,我亲自上手实测了一番,发现两者之间的差异更…

    用户投稿 2026年9月21日
    000
  • Java Stream 高效分组计数并获取Top N元素

    本文深入探讨了如何利用java stream api对数据进行高效的分组计数,并从中提取出现频率最高的top n元素。文章首先介绍了一种简洁的基于全排序的实现方式,该方法适用于数据集较小或top n值接近总数的情况。随后,针对大数据量和小型top n场景下的性能瓶颈,文章详细阐述了如何通过自定义`c…

    2026年9月21日
    000

发表回复

登录后才能评论
关注微信