示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载导读本文基于 flink-learning 仓库中的 flink-learning-connectors-redis 模块完整讲解如何利用 Flink 官方自带的 Redis Connectorflink-connector-redis构建一条「Kafka → Flink → Redis」的实时数据链路先从 Kafka 消费商品事件再经反序列化与字段抽取最终通过RedisSink写入 Redis。文章将覆盖单机 Redis、Redis 集群、Redis Sentinel 三种部署形态的接入配置并结合仓库源码深入RedisMapper、FlinkJedisPoolConfig等核心类的实现原理同时给出配套的数据生产者工具与写入结果验证方式让读者可以直接在本地复现整条链路。一、模块定位Flink 官方自带 Redis Connectorflink-connector-redis是 Flink 官方提供的 Redis 连接器官方版本为 1.1.5它基于 Jedis 客户端实现提供了开箱即用的RedisSink无需自行编写 Redis 客户端连接代码。该模块在仓库中的依赖声明如下pom.xmldependencies dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-redis_2.10/artifactId version1.1.5/version /dependency dependency groupIdredis.clients/groupId artifactIdjedis/artifactId version2.9.0/version /dependency /dependencies其中flink-connector-redis_2.10是 Flink 自带的 Redis 连接器Scala 2.10 版本而redis.clients:jedis是其底层依赖的 Jedis 客户端库RedisSink的所有读写操作最终都委托给 Jedis 完成。适用前提本模块基于flink-connector-redis1.1.5 版本它面向 Flink 1.x 的 DataStream APISource/Sink 接口时代若使用 Flink 1.12 之后的新版本官方已推荐使用flink-connector-redis新制品或直接使用 Table API / SQL 的 Redis 维表方案但本仓库示例所演示的RedisSinkRedisMapper编程模型至今仍是理解 Redis 连接器原理的最佳入门范式。二、整体数据链路与工程结构2.1 数据流向本模块演示的完整链路为Kafka Topiczhisheng └─ FlinkKafkaConsumer 消费SimpleStringSchema 原始 JSON 字符串 └─ Gson 反序列化为 ProductEvent └─ FlatMap 抽取 (id, price) 二元组 └─ RedisSink RedisMapperHSET 写入 Hash key zhisheng └─ Redis单机 / 集群 / Sentinel仓库中 Main.java 是这条链路的完整实现其类注释明确说明「从 Kafka 中读取数据然后写入到 Redis」。2.2 模块文件清单文件作用Main.javaFlink 作业入口消费 Kafka 并写入 RedisProductUtil.javaKafka 生产者工具向 topic 写入商品测试数据RedisTest.java验证类直接读取 Redis 中已写入的数据application.properties作业参数配置Kafka、Redis、并行度等pom.xmlMaven 依赖与打包配置三、核心实现从 Kafka 消费并写入 Redis3.1 作业入口与执行环境Main.java 首先获取StreamExecutionEnvironment与全局参数ParameterToolfinal StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); ParameterTool parameterTool ExecutionEnvUtil.PARAMETER_TOOL; Properties props KafkaConfigUtil.buildKafkaProps(parameterTool);ExecutionEnvUtil.PARAMETER_TOOL是 ExecutionEnvUtil 中的静态参数工具它按照「application.properties文件 → 系统属性 → 命令行参数」的优先级合并参数命令行参数优先级最高。KafkaConfigUtil.buildKafkaProps(parameterTool)由 KafkaConfigUtil 提供它会把application.properties中的kafka.brokers、kafka.zookeeper.connect、kafka.group.id等键写入 Kafka Consumer 的Properties并固定使用StringDeserializer与auto.offset.resetlatest。3.2 Kafka Source消费商品事件作业通过FlinkKafkaConsumer订阅 topic并使用SimpleStringSchema先按原始字符串消费再交给 Gson 反序列化为ProductEvent对象Main.javaSingleOutputStreamOperatorTuple2String, String product env.addSource(new FlinkKafkaConsumer( parameterTool.get(METRICS_TOPIC), // 这个 kafka topic 需要和上面的工具类的 topic 一致 new SimpleStringSchema(), props)) .map(string - GsonUtil.fromJson(string, ProductEvent.class)) // 反序列化 JSON .flatMap(new FlatMapFunctionProductEvent, Tuple2String, String() { Override public void flatMap(ProductEvent value, CollectorTuple2String, String out) throws Exception { // 收集商品 id 和 price 两个属性 out.collect(new Tuple2(value.getId().toString(), value.getPrice().toString())); } });其中topic 名取自配置键metrics.topicPropertiesConstants本模块的 application.properties 中配置为zhishengGsonUtil.fromJson是 GsonUtil 提供的 JSON 反序列化工具ProductEvent是 flink-learning-common 中的商品事件模型Lombok 注解字段包括id、categoryId、code、shopId、name、price以分为单位等flatMap阶段只保留(商品 id, 商品价格)两个字段组成Tuple2String, String作为后续写入 Redis 的数据载体。3.3 Redis Sink通过 RedisMapper 指定写入命令写入 Redis 的核心代码如下Main.java// 单个 Redis FlinkJedisPoolConfig conf new FlinkJedisPoolConfig.Builder().setHost(parameterTool.get(redis.host)).build(); product.addSink(new RedisSinkTuple2String, String(conf, new RedisSinkMapper()));RedisSink构造函数接收两个参数Jedis 连接配置FlinkJedisPoolConfig/FlinkJedisClusterConfig/FlinkJedisSentinelConfig决定连接哪种部署形态的 Redis映射器RedisMapperT决定每条数据以什么 Redis 命令、什么 key 写入。RedisSinkMapper是定义在 Main.java 内部的静态类实现了RedisMapperTuple2String, String接口public static class RedisSinkMapper implements RedisMapperTuple2String, String { Override public RedisCommandDescription getCommandDescription() { return new RedisCommandDescription(RedisCommand.HSET, zhisheng); } Override public String getKeyFromData(Tuple2String, String data) { return data.f0; } Override public String getValueFromData(Tuple2String, String data) { return data.f1; } }RedisMapper接口的三个方法含义如下方法作用getCommandDescription()声明要执行的 Redis 命令及额外信息。这里使用RedisCommand.HSET并指定附加的 Hash 表名zhisheng即写入名为zhisheng的 HashgetKeyFromData(T data)返回写入时使用的 fieldHash 的 field这里取Tuple2的第一个元素商品 idgetValueFromData(T data)返回写入时使用的 value这里取第二个元素商品价格因此最终的效果是对每条商品事件执行HSET zhisheng 商品id 商品价格所有商品的价格都以「id → price」的形式聚合在同一个 Hash keyzhisheng中。RedisCommandDescription支持的命令远不止HSET一种。从 Flink 官方连接器的实现来看RedisCommand枚举覆盖了字符串类SET、GET、INCR、INCRBY、DECR、DECRBY列表类LPUSH、RPUSH、LPOP、RPOP、LLENSet / ZSet 类SADD、SMEMBERS、ZADD、ZREM、ZRANGEHash 类HSET、HGET、HGETALL、HDEL、HLENKey 操作DEL、EXPIRE、TTL、KEYS、TYPE、PING。其中部分命令需要附加参数如HSET需要 Hash 表名、EXPIRE需要过期秒数这正是RedisCommandDescription的第二个参数附加信息存在的意义。读者可根据业务需要替换命令例如统计商品点击量时可使用RedisCommand.INCR。四、三种 Redis 部署形态的接入方式原 README 明确指出 Redis 分三种情况单机 Redis、Redis 集群、Redis Sentinels。Main.java 中三种配置均已给出示例后两种在源码中处于注释状态可取消注释使用4.1 单机 RedisFlinkJedisPoolConfigFlinkJedisPoolConfig conf new FlinkJedisPoolConfig.Builder() .setHost(parameterTool.get(redis.host)) .build(); product.addSink(new RedisSinkTuple2String, String(conf, new RedisSinkMapper()));对应 Jedis 的JedisPool连接池模式redis.host从配置读取本模块默认配置为127.0.0.1见 application.propertiesBuilder还支持.setPort(6379)、.setPassword(...)、.setDatabase(...)、.setTimeout(...)等可选参数未显式设置时使用默认端口6379、默认超时等。4.2 Redis 集群FlinkJedisClusterConfigFlinkJedisClusterConfig clusterConfig new FlinkJedisClusterConfig.Builder() .setNodes(new HashSetInetSocketAddress( Arrays.asList(new InetSocketAddress(redis1, 6379)))).build();对应 Jedis 的JedisCluster用于 Redis Cluster集群分片模式通过setNodes传入集群节点列表InetSocketAddress集合客户端会自动感知集群的 slot 分布与节点拓扑生产环境中节点地址应从配置中心或配置文件读取避免硬编码源码注释也强调「Redis 的 ip 信息一般都从配置文件取出来」。4.3 Redis SentinelFlinkJedisSentinelConfigFlinkJedisSentinelConfig sentinelConfig new FlinkJedisSentinelConfig.Builder() .setMasterName(master) .setSentinels(new HashSet(Arrays.asList(sentinel1, sentinel2))) .setPassword() .setDatabase(1).build();对应 Jedis 的JedisSentinelPool用于 Redis Sentinel哨兵高可用模式setMasterName(master)指定哨兵监控的主节点名称setSentinels(...)传入哨兵节点地址集合客户端通过与哨兵通信获取当前 master 地址并在故障切换后自动感知新的 master可选.setPassword()主节点密码与.setDatabase(1)选择第 1 个逻辑库默认库为 0。三种配置类分别对应 Jedis 三种连接池模型RedisSink在内部会根据配置类型创建对应的 Jedis 实例。实际部署时按自身 Redis 架构选择其一即可其余代码RedisMapper、Sink 注册完全复用。五、配套工具数据生产与写入验证5.1 Kafka 生产者工具 ProductUtil为了让 Flink 作业有数据可消费仓库提供了 ProductUtil.java它作为独立的 Kafka 生产者向 topic 发送商品事件public static final String broker_list localhost:9092; public static final String topic zhisheng; // kafka topic 需要和 flink 程序用同一个 topic for (int i 1; i 10000; i) { ProductEvent product ProductEvent.builder() .id((long) i) // 商品的 id .name(product i) // 商品 name .price(random.nextLong() / 10000000000000L) // 商品价格以分为单位 .code(code i) // 商品编码 .build(); ProducerRecord record new ProducerRecordString, String(topic, null, null, GsonUtil.toJson(product)); producer.send(record); System.out.println(发送数据: GsonUtil.toJson(product)); } producer.flush();关键点生产者与 Flink 作业必须使用同一个 topic代码注释明确提示kafka topic 需要和 flink 程序用同一个 topic即zhisheng消息体使用GsonUtil.toJson序列化为 JSONFlink 端再通过GsonUtil.fromJson反序列化为ProductEvent两者正好构成 JSON 序列化/反序列化的闭环该工具默认向localhost:9092发送 10000 条商品数据可用于本地联调。5.2 写入结果验证 RedisTest数据写入 Redis 后可用 RedisTest.java 直接验证Jedis jedis new Jedis(127.0.0.1); System.out.println(Server is running: jedis.ping()); System.out.println(result: jedis.hgetAll(zhisheng));先通过jedis.ping()确认 Redis 服务连通再通过jedis.hgetAll(zhisheng)读取 Hash keyzhisheng的全部 field-value 对检查 Flink 作业写入的(商品id, 商品价格)数据是否存在该 Hash 表名与RedisSinkMapper.getCommandDescription()中指定的zhisheng严格对应验证了整条链路的正确性。六、运行环境配置与参数说明作业的全部外部参数集中在 application.propertieskafka.brokerslocalhost:9092 kafka.group.idzhisheng kafka.zookeeper.connectlocalhost:2181 metrics.topiczhisheng stream.parallelism4 stream.sink.parallelism4 stream.default.parallelism4 stream.checkpoint.interval1000 stream.checkpoint.enablefalse redis.host127.0.0.1各参数含义与对应源码如下配置键默认值含义消费方kafka.brokerslocalhost:9092Kafka broker 地址KafkaConfigUtilkafka.zookeeper.connectlocalhost:2181Kafka 依赖的 Zookeeper 地址KafkaConfigUtilkafka.group.idzhisheng消费者组 IDKafkaConfigUtilmetrics.topiczhisheng要消费的 Kafka topicMain.javastream.parallelism4作业并行度ExecutionEnvUtilstream.checkpoint.enablefalse是否开启 CheckpointExecutionEnvUtilredis.host127.0.0.1单机 Redis 主机地址Main.java参数加载由ExecutionEnvUtil完成优先级为命令行参数 系统属性 application.properties。因此本地运行时如果 Redis 不在本机可以在提交作业时通过--redis.host ip覆盖配置。七、本地复现步骤基于仓库现有代码在本地完整跑通「Kafka → Flink → Redis」链路的步骤如下准备环境本地启动 Kafka默认localhost:9092与 Redis默认127.0.0.1:6379无密码。启动生产者运行ProductUtil.main()向 topiczhisheng发送 10000 条商品 JSON 数据确保 topic 已创建或让 Kafka 自动创建。启动 Flink 作业运行 Main.java 的main方法作业会从metrics.topic指定的 topic 消费并写入 Redis。若 Redis 不在本机通过--redis.host ip覆盖配置。验证结果运行RedisTest.main()若输出result:{...}中包含形如{商品id商品价格}的键值对即说明链路已打通。需要说明的局限本示例消费 Kafka 采用latest偏移策略KafkaConfigUtil且作业默认关闭 Checkpointstream.checkpoint.enablefalse因此为验证结果建议先启动 Flink 作业、再启动生产者或改为earliest重放数据。八、小结本文以 flink-learning 仓库的 flink-learning-connectors-redis 模块为主线完整还原了「Kafka 商品事件 → Gson 反序列化 → Tuple2 抽取 → RedisSink 写入」的实时链路并分别给出单机、集群、Sentinel 三种 Redis 部署形态的接入代码配合ProductUtil生产端与RedisTest验证端形成了一套可本地复现的端到端示例。从原理层面看RedisSinkRedisMapper这套编程模型把「Redis 命令声明」与「数据字段映射」彻底解耦getCommandDescription()决定写什么命令与命令附加参数getKeyFromData/getValueFromData决定每条数据如何映射到命令参数。掌握这一模型后读者可以自由扩展出 SET、INCR、LPUSH、ZADD 等任意 Redis 写入场景这也是 Flink 连接器设计「Sink Mapper」组合模式的典型代表。延伸阅读仓库中与本文相关的上下游模块还包括 flink-learning-connectors-kafkaKafka 连接器更多用法、flink-learning-commonProductEvent、KafkaConfigUtil、GsonUtil等公共组件以及 flink-learning-project 中的大型项目实战案例。赞分享示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载相关推荐FGO自动刷本助手FGO-py免配置跨平台启动即睡觉养肝护发不是梦FGO自动刷本助手FGO py免配置跨平台启动即睡觉养肝护发不是梦 凌晨一点你盯着手机屏幕上第几百次重复的宝具动画手指机械地滑动。无限池还剩两百池没抽GUI 自动化桌面应用计算机视觉RPA任务调度掌握Ansys ACT二次开发解锁仿真自动化新境界掌握Ansys ACT二次开发解锁仿真自动化新境界 在当今工程仿真领域效率与定制化能力已成为企业竞争力的关键。Ansys ACTAnsys CustomiFlink-Connector-Redis 使用指南Flink Connector Redis 使用指南 项目简介 Flink Connector Redis 是一个基于 Lettuce 的异步 Flink 连接后端消息队列数据集成上一篇GeoJSON.io 实操手册把坐标表格变成能分享的地图全程只需一个浏览器下一篇PaddleSpeech LibriSpeech ASR1 实验基准解读Conformer / Streaming Conformer / Transformer 的 WER 结果与复现指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考