在Apache Flink中读取带键Kafka记录的教程

在apache flink中读取带键kafka记录的教程

本文详细阐述了如何在Apache Flink中使用`KafkaSource`有效读取带键(keyed)的Kafka记录。通过实现自定义的`KafkaRecordDeserializationSchema`,用户可以从Kafka的`ConsumerRecord`中灵活地提取并处理键、值、时间戳、主题、分区及偏移量等元数据,从而克服`valueOnly`反序列化器的局限性,实现更精细的数据处理。

1. 引言:理解带键Kafka记录及其在Flink中的挑战

Apache Kafka作为分布式流处理平台,其消息通常包含一个键(key)和一个值(value)。键在许多场景下至关重要,例如用于消息的有序性保证、特定分区的路由或数据聚合。当使用Kafka控制台生产者工具创建带键记录时,例如:

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

这会生成形如 key:value 的消息,其中 key 和 value 将被Kafka独立处理。

然而,在Apache Flink中,默认的 KafkaSource 配置,特别是使用 KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class) 时,只能获取Kafka记录的值部分,而忽略了键、时间戳以及其他重要的元数据。对于需要根据键进行业务逻辑处理的场景,这显然无法满足需求。本文旨在提供一个全面的教程,指导如何在Flink中正确读取并访问这些带键Kafka记录的所有组成部分。

2. 核心解决方案:自定义 KafkaRecordDeserializationSchema

要在Flink中访问Kafka记录的键、值、时间戳以及其他元数据,核心在于实现一个自定义的 KafkaRecordDeserializationSchema。这个接口允许用户完全控制如何将Kafka的原始 ConsumerRecord 转换成Flink数据流中的元素类型。

KafkaRecordDeserializationSchema 接口定义了几个关键方法,其中最重要的是 deserialize 方法。当Flink从Kafka拉取一条记录时,它会将原始的 ConsumerRecord 传递给此方法。在这个方法内部,我们可以访问 ConsumerRecord 的所有属性,包括:

record.key():获取记录的键。record.value():获取记录的值。record.timestamp():获取记录的时间戳。record.topic():获取记录所属的主题。record.partition():获取记录所在的分区。record.offset():获取记录在分区中的偏移量。record.headers():获取记录的头部信息。

通过自定义 deserialize 方法的逻辑,我们可以将这些信息组合成任何我们需要的输出类型,例如 Tuple2(用于键值对)、自定义的POJO(Plain Old Java Object)或任何其他复杂的数据结构。

wifi优化大师app v1.0.1 安卓版 wifi优化大师app v1.0.1 安卓版

Wifi优化大师最新版是一款免费的手机应用程序,专为优化 Wi-Fi 体验而设计。它提供以下功能:增强信号:提高 Wi-Fi 信号强度,防止网络中断。加速 Wi-Fi:提升上网速度,带来更流畅的体验。Wi-Fi 安检:检测同时在线设备,防止蹭网。硬件加速:优化硬件传输性能,提升连接效率。网速测试:实时监控网络速度,轻松获取网络状态。Wifi优化大师还支持一键连接、密码记录和上网安全测试,为用户提供全面的 Wi-Fi 管理体验。

wifi优化大师app v1.0.1 安卓版 0 查看详情 wifi优化大师app v1.0.1 安卓版

3. 实现自定义反序列化器

以下是一个实现 KafkaRecordDeserializationSchema 的示例,它将Kafka的键和值都反序列化为字符串,并以 Tuple2 的形式输出。

import org.apache.flink.api.common.serialization.DeserializationSchema;import org.apache.flink.api.common.typeinfo.TypeInformation;import org.apache.flink.api.java.tuple.Tuple2;import org.apache.flink.util.Collector;import org.apache.kafka.clients.consumer.ConsumerRecord;import org.apache.kafka.common.serialization.StringDeserializer;import java.io.IOException;/** * 自定义Kafka记录反序列化器,用于提取键和值。 */public class KeyedKafkaRecordDeserializationSchema        implements DeserializationSchema<Tuple2> {    private transient StringDeserializer keyDeserializer;    private transient StringDeserializer valueDeserializer;    @Override    public void open(DeserializationSchema.InitializationContext context) throws Exception {        // 在反序列化器初始化时创建Kafka内置的反序列化器实例        keyDeserializer = new StringDeserializer();        valueDeserializer = new StringDeserializer();    }    @Override    public void deserialize(ConsumerRecord record, Collector<Tuple2> out) throws IOException {        // 从ConsumerRecord中反序列化键和值        String key = keyDeserializer.deserialize(record.topic(), record.headers(), record.key());        String value = valueDeserializer.deserialize(record.topic(), record.headers(), record.value());        // 将键和值作为Tuple2发射出去        out.collect(new Tuple2(key, value));        // 如果需要,还可以访问其他元数据,例如:        // long timestamp = record.timestamp();        // String topic = record.topic();        // int partition = record.partition();        // long offset = record.offset();        // System.out.println("Key: " + key + ", Value: " + value + ", Timestamp: " + timestamp);    }    @Override    public boolean is    EndOfStream(Tuple2 nextElement) {        return false; // 对于流处理,通常返回false    }    @Override    public TypeInformation<Tuple2> getProducedType() {        // 声明此反序列化器生产的类型        return TypeInformation.of(new org.apache.flink.api.java.typeutils.TypeHint<Tuple2>() {});    }}

在上面的示例中:

我们实现了 DeserializationSchema<Tuple2>,表明输出类型是 Tuple2。在 open 方法中,我们初始化了 Kafka 内置的 StringDeserializer 来处理字节数组。deserialize 方法接收原始的 ConsumerRecord,并使用 keyDeserializer 和 valueDeserializer 将字节数组转换为字符串。最后,通过 out.collect(new Tuple2(key, value)) 将反序列化后的键值对发射到 Flink 数据流中。

4. 将自定义反序列化器集成到 KafkaSource

一旦自定义反序列化器 KeyedKafkaRecordDeserializationSchema 实现完成,就可以将其集成到 KafkaSource 的构建过程中。只需调用 setDeserializer() 方法,传入自定义反序列化器的实例。

import org.apache.flink.connector.kafka.source.KafkaSource;import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.api.java.tuple.Tuple2;import org.apache.flink.api.common.eventtime.WatermarkStrategy;public class FlinkKeyedKafkaConsumer {    public static void main(String[] args) throws Exception {        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();        String bootstrapServers = "localhost:9092"; // 替换为你的Kafka服务器地址        String topic = "test3";        String groupId = "my-flink-consumer-group";        // 构建KafkaSource,使用自定义的反序列化器        KafkaSource<Tuple2> source = KafkaSource.<Tuple2>builder()                .setBootstrapServers(bootstrapServers)                .setTopics(topic)                .setGroupId(groupId)                .setStartingOffsets(OffsetsInitializer.earliest()) // 从最早的偏移量开始消费                .setDeserializer(new KeyedKafkaRecordDeserializationSchema()) // 使用自定义反序列化器                .build();        // 从Kafka源创建数据流        DataStream<Tuple2> stream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Keyed Source");        // 对数据流进行处理,例如打印键和值        stream.map(record -> "Receiving from Kafka - Key: " + record.f0 + ", Value: " + record.f1)              .print();        // 执行Flink作业        env.execute("Flink Keyed Kafka Consumer Job");    }}

在上述代码中,我们用 new KeyedKafkaRecordDeserializationSchema() 替换了之前的 KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class)。现在,stream 中的元素类型将是 Tuple2,其中 f0 代表键,f1 代表值。

5. 注意事项与最佳实践

处理 null 键或值: Kafka记录的键或值可能为 null。在自定义反序列化器中,record.key() 或 record.value() 返回的字节数组也可能为 null。在进行反序列化时,务必检查这些 byte[] 是否为 null,以避免 NullPointerException。例如:

String key = (record.key() == null) ? null : keyDeserializer.deserialize(record.topic(), record.headers(), record.key());String value = (record.value() == null) ? null : valueDeserializer.deserialize(record.topic(), record.headers(), record.value());

选择合适的输出数据类型:Tuple2: 简单直接,适用于键值都是基本类型的情况。自定义POJO: 当需要同时访问键、值、时间戳、主题、分区、偏移量等多个属性,并且这些属性共同构成一个逻辑实体时,自定义POJO是更好的选择。例如:

public class MyKafkaRecord {    public String key;    public String value;    public long timestamp;    public String topic;    // ... 其他字段}// 在 deserialize 方法中创建并填充 MyKafkaRecord 实例

使用POJO时,确保POJO类有公共的无参构造函数,并且所有字段都是公共的或有公共的getter/setter方法,以便Flink能够正确地进行类型序列化和反序列化。

Row 类型: Flink SQL和Table API通常使用 Row 类型。如果你的数据流最终会与Table API/SQL集成,考虑将输出类型设计为 Row。错误处理: 反序列化过程中可能会发生错误,例如数据格式不匹配。在 deserialize 方法中,可以捕获 IOException 或其他相关的异常。根据业务需求,可以选择跳过错误记录、记录日志、将错误记录发送到旁路输出(side output)或使作业失败。性能考量: 自定义反序列化器会在每个记录上执行。确保 deserialize 方法的逻辑尽可能高效。避免在 deserialize 方法中执行耗时的操作,例如网络请求或复杂的计算。类型信息: getProducedType() 方法必须准确返回反序列化器产生的类型信息。这对于Flink的类型系统至关重要。对于复杂类型(如POJO或泛型类型),可以使用 TypeInformation.of(new TypeHint() {}) 来获取正确的类型信息。

6. 总结

通过实现自定义的 KafkaRecordDeserializationSchema,Apache Flink用户可以完全控制如何从Kafka的 ConsumerRecord 中提取和处理数据,包括键、值、时间戳以及其他元数据。这种方法提供了极大的灵活性,使得Flink能够充分利用Kafka带键记录的丰富语义,从而构建出更强大、更精细的流处理应用。理解并正确运用自定义反序列化器是开发高效、健壮的Flink Kafka集成应用的关键一步。

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

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

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
10期实战直播|腾讯云可观测平台全面升级,场景实践一次讲透
上一篇 2025年12月1日 21:26:53
2024Q2全球入门手机TOP10出炉:Redmi 13C屠榜第一 遥遥领先
下一篇 2025年12月1日 21:26:59

相关推荐

  • ​​VSCode的进阶玩法!这些快捷键组合让你秒变编程大神​​

    要高效利用vscode进行代码导航,1. 使用ctrl + p / cmd + p快速打开文件;2. 使用ctrl + shift + o / cmd + shift + o定位文件中的符号;3. 使用ctrl + – / cmd + -后退,ctrl + shift + –…

    2026年9月24日
    700
  • AI PC 新晋狠角色:5000 元价位 Arrow Lake 最优解 惠普战 66 2025 争当全能卷王

    AI PC 新晋狠角色:5000 元价位 Arrow Lake 最优解 惠普战 66 2025 争当全能卷王AI PC 新晋狠角色:5000 元价位 Arrow Lake 最优解 惠普战 66 2025 争当全能卷王AI PC 新晋狠角色:5000 元价位 Arrow Lake 最优解 惠普战 66 2025 争当全能卷王AI PC 新晋狠角色:5000 元价位 Arrow Lake 最优解 惠普战 66 2025 争当全能卷王

    在 5000 元级别的主流商务本市场,长久以来似乎都遵循着一套 ” 潜规则 “:追求性能就得牺牲便携,看重耐用又往往在外观和屏幕上妥协,想要全面的接口以及优质的售后服务,预算就得一加再加。但现在,一个 ” 新晋狠角色 ” 决意打破这一局面。 惠普商用产…

    2026年9月24日 用户投稿
    000
  • 智能平权下,燃油车如何升级?

    智能平权下,燃油车如何升级?智能平权下,燃油车如何升级?智能平权下,燃油车如何升级?智能平权下,燃油车如何升级?

    曾几何时,“智能驾驶是电动车的专属”成为汽车行业的共识。宝马、奔驰、奥迪等传统豪华品牌长期专注于机械精密性和驾驶质感,在智能化布局上尤为谨慎,一度被贴上保守与落后的标签。 与此同时,新能源品牌凭借智能化迅速打开市场缺口,成功构建起“电动即智能、燃油即传统”的认知框架,在舆论和市场销量中占据先机。 ☞…

    2026年9月24日 用户投稿
    000
  • mac怎么使用听写功能_mac听写输入开启方法

    首先启用高级听写功能,进入系统设置→键盘→听写,勾选“使用高级听写”并下载语言包;随后可设置快捷键(如双击Fn键)快速启动语音输入;在支持的应用中也可通过菜单栏“编辑→开始听写”直接调用;最后根据需要配置听写语言、自动纠正及连续听写选项以提升识别准确率。 如果您希望在Mac上通过语音输入文字以提高效…

    2026年9月24日
    200
  • WPS怎么办批量打印文件_WPS多文档批量打印与队列管理步骤

    使用WPS内置批量打印功能,通过工具→批量处理→批量打印添加文件并设置参数后一键打印;2. 在资源管理器中多选文件右键以WPS打开,再在文件菜单中选择打印所有文档;3. 通过控制面板的设备和打印机进入打印队列,可实时监控、暂停或取消任务,确保批量打印过程稳定可控。 如果您需要在短时间内处理多个WPS…

    2026年9月24日
    400
  • Intel OpenCAS缓存加速方案

    open cas 架构概览:数据从hdd盘读取后被复制到open cas的缓存中,后续的读取操作从内存中进行,从而提高读写效率。在write-through模式下,所有数据同步刷新到open cas的ssd和后端的hdd中。在write-back模式下,数据同步写入到open cas的ssd中,然后…

    2026年9月24日
    400
  • 使用MySQL命令行客户端进行交互式管理

    使用MySQL命令行客户端进行交互式管理使用MySQL命令行客户端进行交互式管理使用MySQL命令行客户端进行交互式管理使用MySQL命令行客户端进行交互式管理

    mysql命令行客户端的常用命令包括:1. 使用mysql -u 用户名 -p命令连接数据库;2. 执行show databases;查看所有数据库;3. 使用use 数据库名;选择数据库;4. 使用select * from 表名;查询数据;5. 使用insert into 表名 (列1, 列2)…

    2026年9月24日 用户投稿
    500
  • windows8如何取消开机密码_windows8取消开机密码的操作

    1、通过netplwiz取消密码验证并设置自动登录;2、在控制面板中更改密码为空实现无密码开机;3、使用命令提示符执行net user命令清除账户密码,三者均可实现Windows 8.1启动时跳过密码输入。 如果您希望在启动Windows 8系统时跳过手动输入密码的步骤,可以直接通过系统内置工具配置…

    2026年9月24日
    000
  • 赋能AI未来!康盈半导体 AI 应用存储新品登陆 elexcon 2025 展会

    赋能AI未来!康盈半导体 AI 应用存储新品登陆 elexcon 2025 展会赋能AI未来!康盈半导体 AI 应用存储新品登陆 elexcon 2025 展会赋能AI未来!康盈半导体 AI 应用存储新品登陆 elexcon 2025 展会赋能AI未来!康盈半导体 AI 应用存储新品登陆 elexcon 2025 展会

    8 月 26 日,中国电子、嵌入式及半导体先进封测行业的风向标 ——elexcon2025 深圳国际电子展暨嵌入式展盛大开幕。作为本届展会的重磅环节之一,国产存储领军品牌康盈半导体携新而来,以 “小而不凡,速启 ai 未来” 为核心主题,正式发布 2025 年存储新品,同步拉开面向 ai 终端应用的…

    2026年9月24日 用户投稿
    000
  • Java中固定长度用户ID输入验证:解决int类型长度检查问题

    本文详细介绍了在Java程序中如何实现用户输入固定长度ID的验证机制。针对常见的int cannot be dereferenced错误,我们将探讨将ID作为字符串读取并进行长度及格式校验的最佳实践,并提供处理字母数字型和纯数字型ID的示例代码,确保数据输入的准确性和程序的健壮性。 引言:用户输入验…

    2026年9月24日
    500
  • 手机qq浏览器阅读模式怎么退出_手机QQ浏览器退出阅读模式操作方法

    要退出手机QQ浏览器的阅读模式,首先点击页面顶部的“阅读模式已开启”按钮即可立即退出;若无明显开关,可刷新页面或长按刷新强制重新加载;也可通过右上角“更多”菜单选择“退出阅读模式”;为避免再次进入,可在设置中关闭“自动进入阅读模式”功能。 如果您在使用手机QQ浏览器时进入了阅读模式,但希望恢复到原始…

    2026年9月24日
    200
  • 数据实时迁移同步工具 CloudCanal v5.2.0.0 发布,支持 SaaS 全托管

    cloudcanal 免费社区版 是 clougence 公司推出的一款全自研、可视化、自动化数据迁移同步工具,具备 结构迁移、数据迁移、数据同步、数据校验、数据订正 等功能,支持 60+ 款流行关系型数据库、实时数仓、消息中间件、缓存数据库和搜索引擎之间数据互通,其中包含国产数据库 oceanba…

    2026年9月24日
    000
  • VSCode如何实现AI代码反混淆 VSCode智能分析混淆代码的技巧

    vscode没有一键ai反混淆功能,但可通过智能扩展、调试器、ast查看器、代码格式化工具及外部ai工具集成来辅助分析和逐步还原混淆代码;2. 利用eslint、prettier等扩展提升代码可读性,通过“重命名符号”“转到定义”“查找引用”等功能追踪变量和函数流向,结合多光标编辑和代码片段进行手动…

    2026年9月24日
    100
  • Laravel 表单验证失败后保留输入值:最佳实践教程

    本文旨在帮助 Laravel 开发者解决表单验证失败后,如何保留用户已输入数据的问题。我们将深入探讨 withInput() 方法的使用,并提供清晰的代码示例,确保即使在验证失败的情况下,用户体验也能保持流畅。通过本文的学习,你将掌握在 Laravel 中优雅地处理表单验证,并提升应用的可用性。 在…

    2026年9月24日
    000
  • 小红书推广选择阅读量还是粉丝量?小红书怎么推广引流

    小红书作为融合内容、社交与电商的综合性平台,近年来吸引了大量创作者和品牌入驻。在进行推广时,很多人常常纠结:是更重视阅读量,还是更关注粉丝量?本文将从两者的定义出发,分析各自的优劣势,并提供实用建议,帮助你制定适合自己的推广策略。 一、阅读量与粉丝量的本质区别 1. 阅读量 阅读量代表的是某篇笔记或…

    2026年9月24日
    000
  • 怎么在mysql中创建数据库表 mysql建表完整流程解析

    在 mysql 中创建数据库表的步骤包括:1) 选择合适的数据类型,如 int、varchar、timestamp;2) 设置索引,如主键和唯一索引;3) 应用约束条件,如 not null 和 unique;4) 设计表结构以满足业务需求,如使用 foreign key 和 enum;5) 优化性…

    2026年9月24日
    000
  • 电脑网络连接限制为100Mbps怎么办_网速被限制的解决方法

    首先确认网卡是否支持千兆速率,进入设备管理器查看网络适配器型号并核实规格;接着在网卡属性中将“速度和双工”设置为1.0 Gbps全双工或自动协商;更换符合Cat 5e及以上标准的网线,确保物理连接可靠;检查路由器或交换机端口设置,确保未强制限制为100Mbps;更新或重装网卡驱动程序至最新版本;最后…

    2026年9月24日
    100
  • 如何通过日志排查权限问题

    排查权限问题需从日志入手,重点分析时间、用户、资源路径、拒绝原因及调用堆栈。首先检查应用日志中“用户无权访问”等提示,结合Web服务器日志中的403/401状态码定位请求异常;再查看操作系统日志如/var/log/secure中SSH或sudo拒绝记录,确认系统级权限问题;同时审查中间件如Sprin…

    2026年9月24日
    100
  • win10软件不兼容怎么办_win10软件兼容性处理方法

    首先使用兼容性疑难解答工具检测并修复问题,若无效则手动设置兼容模式为Windows 7或8,同时安装必要的Visual C++和.NET运行库,更新显卡等驱动程序,并尝试以管理员身份运行程序。 如果您尝试在Windows 10系统上运行某个软件,但出现“此应用无法在你的电脑上运行”或程序闪退等错误提…

    2026年9月24日
    000
  • 生成Java中全范围正Double随机数的正确方法

    本文旨在指导开发者如何在Java中生成覆盖整个正Double范围的随机数,并解释了使用ThreadLocalRandom.nextDouble(Double.MIN_VALUE, Double.MAX_VALUE)可能产生偏差的原因。我们将提供一种基于位操作的替代方案,确保生成的随机数在Double…

    2026年9月24日
    100

发表回复

登录后才能评论
关注微信