ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Flink Connector 实战:Flink 消费 Kafka 商品数据写入 Redis(flink-learning Redis Connector 用法详解)

Flink Connector 实战:Flink 消费 Kafka 商品数据写入 Redis(flink-learning Redis Connector 用法详解) 示例工程大数据【免费下载链接】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点击查看免费下载导读在生产环境中流式计算的结果常常需要落地到 Redis供第三方应用按 key 快速查询——例如根据商品 ID 获取商品价格。本文以 flink-learning 仓库中的 Redis Connector 实战案例为主线完整讲解Kafka 模拟商品数据 → Flink 消费并提取 ID/价格 → 写入 Redis → 验证读取的全链路实现并深入分析 Redis Connector 的三种连接方式单机 / 集群 / Sentinel、RedisMapper接口原理与依赖选型让读者既能直接跑通案例也能理解其底层机制。1. 案例背景与整体架构本案例对应的完整工程位于 flink-learning-connectors/flink-learning-connectors-redis 模块其 README 明确了该模块的定位利用自带的 Redis Connector 从 Kafka 中读取数据然后写入到 Redis。整体数据处理链路如下ProductUtil模拟商品数据 │ 发送 JSON 到 Kafka topic zhisheng ▼ Flink JobFlinkKafkaConsumer 消费 Kafka │ 反序列化为 ProductEvent ▼ flatMap 提取 (商品ID, 商品价格) │ ▼ RedisSinkRedisSinkMapperHSET 命令写入 keyzhisheng 的 Hash │ ▼ 第三方服务 / Jedis 客户端 按商品ID 查询价格核心思路是从 Kafka 读取所有商品信息仅提取商品 ID 与商品价格两个字段写入 Redis供第三方服务根据商品 ID 查询对应价格。2. 安装与启动 Redis2.1 源码方式安装从 Redis 官网下载源码包并编译以 redis-5.0.4 为例wget http://download.redis.io/releases/redis-5.0.4.tar.gz tar xzf redis-5.0.4.tar.gz cd redis-5.0.4 make2.2 HomeBrew 方式安装macOSbrew install redis若需要后台常驻运行 Redis 服务brew services start redis安装完成后可在/usr/local/bin目录下找到两个关键命令redis-server启动服务端redis-cli启动客户端执行redis-server打开服务端后另开一个终端执行redis-cli即可进入交互式客户端后续可用hgetall、hget等命令验证写入结果。3. 模拟商品数据发送到 Kafka3.1 商品事件模型 ProductEvent商品类定义在 flink-learning-common/src/main/java/com/zhisheng/common/model/ProductEvent.java其中id商品 ID与price商品价格以分为单位是本案例最终要提取并写入 Redis 的字段Data Builder AllArgsConstructor NoArgsConstructor public class ProductEvent { /** Product Id */ private Long id; /** Product 类目 Id */ private Long categoryId; /** Product 编码 */ private String code; /** Product 店铺 Id */ private Long shopId; /** Product 店铺 name */ private String shopName; /** Product 品牌 Id */ private Long brandId; /** Product 品牌 name */ private String brandName; /** Product name */ private String name; /** Product 图片地址 */ private String imageUrl; /** Product 状态1(上架),-1(下架),-2(冻结),-3(删除) */ private int status; /** Product 类型 */ private int type; /** Product 标签 */ private ListString tags; /** Product 价格以分为单位 */ private Long price; }该模型使用 Lombok 注解Data、Builder等因此可以像ProductEvent.builder().id(...).name(...).build()这样链式构建对象。3.2 数据发送工具类 ProductUtil仓库中 ProductUtil.java 负责不断向 Kafka 模拟发送商品数据关键实现如下public class ProductUtil { public static final String broker_list localhost:9092; public static final String topic zhisheng; //kafka topic 需要和 flink 程序用同一个 topic public static final Random random new Random(); public static void main(String[] args) { Properties props new Properties(); props.put(bootstrap.servers, broker_list); props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); KafkaProducer producer new KafkaProducerString, String(props); 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(); } }要点说明topic 必须与 Flink 程序消费的 topic 一致这里固定为zhisheng序列化器使用StringSerializer发送的是GsonUtil.toJson(product)序列化后的 JSON 字符串循环发送 10000 条模拟数据后调用producer.flush()强制刷新缓冲确保消息全部送达商品价格通过random.nextLong() / 10000000000000L生成一个随机的以分为单位的价格数值即分与ProductEvent.price的字段注释一致。4. Flink 消费 Kafka 中的商品数据Flink Job 主程序为 flink-learning-connectors-redis 的 Main.java。第一步是消费 Kafka 并提取 id 与 pricepublic class Main { public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); ParameterTool parameterTool ExecutionEnvUtil.PARAMETER_TOOL; Properties props KafkaConfigUtil.buildKafkaProps(parameterTool); SingleOutputStreamOperatorTuple2String, 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())); } }); env.execute(flink redis connector); } }实现细节说明Kafka 配置统一封装KafkaConfigUtil.buildKafkaProps(parameterTool)在 flink-learning-common 的 KafkaConfigUtil.java 中实现它从ParameterTool读取配置并写入bootstrap.servers默认localhost:9092、zookeeper.connect默认localhost:2181、group.id默认zhisheng等属性并固定使用StringDeserializer反序列化 key/value、auto.offset.resetlatesttopic 读取自配置parameterTool.get(METRICS_TOPIC)对应 PropertiesConstants.java 中的常量METRICS_TOPIC metrics.topic即配置文件application.properties中的metrics.topic项需要与ProductUtil发送的 topic 一致反序列化使用GsonUtil.fromJson(string, ProductEvent.class)将 Kafka 中的 JSON 字符串解析为ProductEvent对象字段提取flatMap将每条商品记录转换为Tuple2String, Stringf0为商品 ID 字符串、f1为价格字符串这正是接下来写入 Redis 的 KV 数据形态。在 IDEA 中先启动 Flink Job再运行ProductUtil注意需提前启动 KafkaJob 即可消费到商品的 id 与 price 数据。调试阶段可先通过product.print()将结果打印到控制台验证确认数据正确后再接入 Redis Sink。5. Redis Connector 简介与依赖选型5.1 Redis Connector 支持的三种环境Redis Connector 提供用于向 Redis 发送数据的接收器Sink可支持三种不同类型的 Redis 环境单 Redis 服务器StandaloneRedis 集群ClusterRedis Sentinel哨兵模式对应到源码中正是 Main.java 里注释展示的三类配置类FlinkJedisPoolConfig单机、FlinkJedisClusterConfig集群、FlinkJedisSentinelConfigSentinel。5.2 Maven 依赖官方依赖Apache Flink 老版本官方 Redis Connector 只有较老的版本此后一直未更新dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-redis_2.10/artifactId version1.1.5/version /dependencyBahir 社区依赖在 Apache Bahir 项目下也能找到功能相同的 Flink Redis Connectordependency groupIdorg.apache.bahir/groupId artifactIdflink-connector-redis_2.11/artifactId version1.0/version /dependency两个依赖的功能一致仓库演示选用的是官方flink-connector-redis_2.10。实际工程 pom.xml 中确认了该依赖同时显式引入了jedis 2.9.0作为 Redis Java 客户端dependency groupIdredis.clients/groupId artifactIdjedis/artifactId version2.9.0/version /dependency注意该连接器基于老版本 Flink API如FlinkKafkaConsumer等若在较新版本的 Flink 中使用需关注 API 兼容性并做相应适配。5.3 打包配置pom.xml 中还配置了maven-shade-plugin将com.zhisheng.connectors.redis.Main声明为mainClass因此可通过mvn package打出可直接提交运行的 Fat Jar。6. Flink 写入数据到 Redis6.1 单机 Redis 配置在 Main.java 中单机 Redis 的 Sink 配置与接入代码如下//单个 Redis FlinkJedisPoolConfig conf new FlinkJedisPoolConfig.Builder().setHost(parameterTool.get(redis.host)).build(); product.addSink(new RedisSinkTuple2String, String(conf, new RedisSinkMapper()));FlinkJedisPoolConfig用于描述单机 Redis 的连接池配置这里通过Builder链式设置host从参数redis.host读取实际部署时通常配置在application.properties中而不是硬编码RedisSink是通用的 Redis 输出算子需要两个参数连接配置 RedisMapper实现后者负责描述如何把数据映射为 Redis 命令。6.2 集群与 Sentinel 配置源码注释中的备选方案仓库源码以注释形式给出了另外两种环境的配置模板供读者按需启用Redis 集群FlinkJedisClusterConfig clusterConfig new FlinkJedisClusterConfig.Builder() .setNodes(new HashSetInetSocketAddress( Arrays.asList(new InetSocketAddress(redis1, 6379)))).build();Redis SentinelFlinkJedisSentinelConfig sentinelConfig new FlinkJedisSentinelConfig.Builder() .setMasterName(master) .setSentinels(new HashSet(Arrays.asList(sentinel1, sentinel2))) .setPassword() .setDatabase(1).build();可见三种模式分别对应FlinkJedisPoolConfig/FlinkJedisClusterConfig/FlinkJedisSentinelConfig三个配置类写法高度一致均使用 Builder 模式。6.3 RedisMapper 接口实现写入规则的核心是RedisMapper接口。仓库在 Main.java 内部定义了一个静态内部类RedisSinkMapperpublic 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; } }三个方法各司其职getCommandDescription()声明本次写入使用的 Redis 命令与附加参数。这里使用RedisCommand.HSET命令附加参数zhisheng作为Hash 的 key表名即最终写入的是名为zhisheng的 Hash 结构getKeyFromData(data)返回 Hash 中的field这里取Tuple2的f0商品 ID 字符串getValueFromData(data)返回 Hash 中该 field 对应的value这里取f1商品价格字符串。RedisCommand枚举支持HSET、SET、SADD、LPUSH、INCR等常见 Redis 命令开发者可根据业务选择若使用HSET这类需要附加 key 的命令则在RedisCommandDescription构造函数的第二个参数中传入附加 key。6.4 完整 Sink 接入product.addSink(new RedisSinkTuple2String, String(conf, new RedisSinkMapper())); env.execute(flink redis connector);将flatMap得到的Tuple2String, String数据流接入RedisSink提交作业后每条数据都会以HSET zhisheng {商品ID} {商品价格}的形式写入 Redis。7. 项目运行与验证7.1 运行步骤启动 Kafkabootstrap.serverslocalhost:9092与 Redis 服务端redis-server在application.properties参考 PropertiesConstants.java 中定义的metrics.topic、kafka.brokers等配置项中配置好metrics.topic与发送工具 topic 一致如zhisheng以及redis.host如localhost在 IDEA 中先运行 Main.java 启动 Flink Job再运行 ProductUtil.java 向 Kafka 发送模拟数据观察 Job 控制台输出与 Redis 中的数据。7.2 用 Jedis 客户端验证写入结果仓库在 RedisTest.java 中提供了最直接的验证方式——用 Jedis 客户端连接 Redis 并读取整个 Hashimport redis.clients.jedis.Jedis; public class RedisTest { public static void main(String[] args) { 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()返回PONG说明 Redis 连接正常jedis.hgetAll(zhisheng)返回 Hash 中全部 field-value 对即可看到以商品 ID 为 field、以价格为 value 的完整映射从而验证 Flink → Kafka → Redis 全链路数据已正确落地。也可以直接用redis-cli验证redis-cli hgetall zhisheng hget zhisheng 1 # 查询商品ID1 的价格第三方服务同样可通过HGET zhisheng {商品ID}快速获取商品价格。8. 小结与反思本案例完整演示了 flink-learning 仓库中 Redis Connector 的实战用法全链路闭环从 Kafka 模拟数据源ProductUtil、Flink 消费与字段提取Main中的FlinkKafkaConsumerflatMap到RedisSink写入再到 Jedis 客户端RedisTest读取验证每个环节都有可直接运行的仓库源码支撑三种环境适配Redis Connector 通过FlinkJedisPoolConfig单机、FlinkJedisClusterConfig集群、FlinkJedisSentinelConfigSentinel三种配置类覆盖了常见的 Redis 部署形态Mapper 抽象设计RedisMapper将数据到 Redis 命令的映射逻辑与 Sink 解耦通过getCommandDescription/getKeyFromData/getValueFromData三个方法即可定制任意写入规则这也是 Redis Connector 易用的核心所在。需要提醒的几点依赖版本较老官方flink-connector-redis_2.101.1.5与 Bahir 版本均久未更新接入新版本 Flink 时需评估兼容性必要时基于RedisMapper思路自研或采用新版连接器配置外置redis.host等连接参数应从配置文件读取避免硬编码仓库中ExecutionEnvUtil、ParameterTool的组合使用正是这一思路的体现价格单位ProductEvent.price以分为单位第三方读取时需注意单位换算。更多 Redis Connector 的底层实现细节可继续阅读 flink-learning-connectors-redis 模块源码以及 flink-learning 系列其他 Connector 实战章节如 Kafka、Elasticsearch 等对照学习。赞分享示例工程大数据【免费下载链接】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 自带 Redis Connector 实现 Kafka 到 Redis 的数据写入单机 / 集群 / Sentinel 三模式flink learning 实战Flink 自带 Redis Connector 实现 Kafka 到 Redis 的数据写入单机 / 集群 / Sent示例工程大数据FGO自动刷本助手FGO-py免配置跨平台启动即睡觉养肝护发不是梦FGO自动刷本助手FGO py免配置跨平台启动即睡觉养肝护发不是梦 凌晨一点你盯着手机屏幕上第几百次重复的宝具动画手指机械地滑动。无限池还剩两百池没抽GUI 自动化桌面应用计算机视觉RPA任务调度掌握Ansys ACT二次开发解锁仿真自动化新境界掌握Ansys ACT二次开发解锁仿真自动化新境界 在当今工程仿真领域效率与定制化能力已成为企业竞争力的关键。Ansys ACTAnsys Customi创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表