ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Flink读取Kafka数据:双写Redis与MySQL的实时链路实践

Flink读取Kafka数据:双写Redis与MySQL的实时链路实践 简介面向大数据流计算开发者和实时数仓初学者这是一个围绕 Apache Flink 消费 Kafka 消息、完成窗口聚合后再写入 Redis 集群与 MySQL 的完整工程示例。它覆盖了从数据接入、流式处理到结果存储的典型链路可帮助读者快速理解并复用 Flink 与外部存储的集成方式适合需要搭建实时监控、日志分析、在线指标计算等场景的开发者。压缩包共包含145个文件整体大小约48.47MB主要构成为 Java 源码、XML 工程配置、class 编译产物和 properties 属性配置另含少量 jar 依赖及命令行脚本其中源码与配置便于按模块跟踪逻辑class 产物则方便直接部署或验证。当前已有597人学习。项目具体展示了 Kafka 连接参数与消费组设置、基于 keyBy 和窗口算子的计算过程、Redis Sink 的集群写入方式以及通过 JDBC 将流式结果导入 MySQL 的批量提交策略同时包含集群槽位分配与数据路由相关配置能够为实际工程中常见的多组件协同问题提供排错思路和代码参照。1. flink读取kafka数据一份编译好的实时链路样本拿到的这份flink读取kafka数据.zip表面看只是一堆编译后的 class 文件但把它反推回去其实是一个完整的 Flink 实时链路样本从 Kafka 消费日志事件做窗口计算再双写 Redis 集群和 MySQL。它不是源码教学包而是给你“对照验证”用的——你写的 Flink 作业跑出来的行为跟这份 class 的行为是否一致。适合正在搭实时日志处理、又不想从零画架构的开发者拿它当链路骨架的参照物。我拆完后发现里面水印提取、双 Sink 的参数设置比预期更细值得逐段展开。2. 从 class 文件反推数据模型LogEvent 与 Schema 的结构线索2.1 为什么先看 class 清单而不是直接找代码打开 zip 先别急着找源码这个包里确实没有.java只有编译产物。class 文件名本身就是最好的设计文档。我从清单里提取出几组关键类LogEvent、LogEventSchema、LogEventApp、LogEventWaterMarkExtractorReqInfo、RequestMessage、ResponseMessagePropertiesConfiguration这组命名说明核心事件模型是一个日志事件LogEvent里面有请求信息ReqInfo请求消息和响应消息是它的两个主要构成部分。LogEventSchema负责序列化和反序列化LogEventWaterMarkExtractor负责从事件里提取水位线。PropertiesConfiguration管理外部连接配置包括 Kafka、Redis、MySQL 三套。先把这个结构确认了后面的反推才有根据。2.2 用 javap 反编译查看类签名拿到 class 后我第一个动作是用javap看公开方法和字段签名不需要反编译所有实现看签名就能确认数据模型。javap 是 JDK 自带的工具不需要额外装东西javap -p -c LogEvent.class-p显示私有成员-c打印方法字节码。如果是反编译全部逻辑可以用cfr或fernflower但做架构判断时javap已经够用。LogEvent类里通常会有getReqInfo()、getTimestamp()这类方法看到时间戳字段的 getter就可以确认水印提取器是基于事件时间而不是处理时间。2.3 事件时间与水印提取的关系确认LogEventWaterMarkExtractor这个类名值得停下来多看一眼。Flink 里做窗口聚合最容易被坑的就是时间语义。如果在 class 里看到assignTimestampsAndWatermarks的调用点说明作业用的是事件时间。日志类数据天然带客户端时间戳用事件时间计算出的延迟指标才有业务意义否则跑出来的数字全是服务器处理时刻参考价值大打折扣。这里我还原一下典型的水印提取写法这个也是我在真实项目里常用的模式DataStreamLogEvent withWatermarks stream .assignTimestampsAndWatermarks( new LogEventWaterMarkExtractor() ); public class LogEventWaterMarkExtractor extends BoundedOutOfOrdernessTimestampExtractorLogEvent { public LogEventWaterMarkExtractor() { super(Time.seconds(10)); } Override public long extractTimestamp(LogEvent element) { return element.getTimestamp(); } }BoundedOutOfOrdernessTimestampExtractor的意思是允许乱序数据最多迟到 10 秒超过这个范围的数据会被丢弃。extractTimestamp返回毫秒时间戳。这里有个参数值得记住10秒不是拍脑袋定的要看上游 Kafka 里日志产生时间和到达时间的差值分布。一般先跑一天数据取 P95 的延迟作为允许乱序的阈值不要一开始就设 60 秒那会让窗口计算结果严重滞后。2.4 PropertiesConfiguration三套连接配置的集中管理PropertiesConfiguration是链路里最实用的一环。它把 Kafka、Redis、MySQL 三套配置统一加载到一个 Properties 对象里再分发给各个客户端。这个设计我比较认可因为一个实时作业的配置项少说二十个如果散落在代码里换环境时改到怀疑人生。常见做法是PropertiesConfiguration config new PropertiesConfiguration(); Properties kafkaProps config.getKafkaProperties(); Properties redisProps config.getRedisProperties(); Properties mysqlProps config.getMysqlProperties();getKafkaProperties()内部通常会加载包含bootstrap.servers、group.id、auto.offset.reset的配置getRedisProperties()包含集群节点列表和密码getMysqlProperties()包含 JDBC URL 和用户名密码。后面接 Sink 时直接从这个配置类拿 Properties 对象传给连接器省去一堆重复代码。3. 核心消费链路Kafka Source 的参数设置与反序列化3.1 Kafka Source 初始化从 class 清单来看LogEventApp是作业入口里面必然有addSource创建 Kafka Consumer 的逻辑。Flink 1.x 用FlinkKafkaConsumer新项目建议直接用 Kafka Source但既然这个包里有LogEventSchema我先按兼容写法讲。初始化 Kafka Source 的常见写法和关键参数如下Properties kafkaProps new Properties(); kafkaProps.setProperty(bootstrap.servers, 10.0.0.11:9092,10.0.0.12:9092); kafkaProps.setProperty(group.id, log-event-app); kafkaProps.setProperty(auto.offset.reset, latest); KafkaSourceLogEvent source KafkaSource.LogEventbuilder() .setBootstrapServers(10.0.0.11:9092,10.0.0.12:9092) .setTopics(log-event-topic) .setGroupId(log-event-app) .setStartingOffsets(OffsetResetStrategy.LATEST) .setDeserializer(new LogEventSchema()) .build(); DataStreamLogEvent stream env.fromSource( source, WatermarkStrategy.noWatermarks(), log-event-kafka-source );bootstrap.servers只要写集群中任意几个 broker 地址即可客户端会通过它们发现完整的 broker 列表不需要把全部节点写进去。group.id决定消费者组的归属同一 group 内的消费者会分担不同分区的消费。auto.offset.reset只在当前 group 没有提交过 offset 时生效新 group 设置成latest意味着从头开始只消费新数据而earliest会从最早的数据开始重放。实时日志场景我一般用latest避免作业刚启动就灌入几天的历史数据。3.2 LogEventSchema 的反序列化实现LogEventSchema是链路的第一道关口它把 Kafka 里的字节数组还原成LogEvent对象。用javap看它的方法应该有deserialize和serialize两个方向。反序列化的实现一般是手写 JSON 解析或使用快速 JSON 库这段逻辑要特别注意容错处理。public class LogEventSchema implements DeserializationSchemaLogEvent { Override public LogEvent deserialize(byte[] message) throws IOException { String json new String(message, StandardCharsets.UTF_8); JSONObject obj JSON.parseObject(json); LogEvent event new LogEvent(); event.setReqInfo(obj.getJSONObject(reqInfo).toJavaObject(ReqInfo.class)); event.setRequestMessage(obj.getJSONObject(requestMessage).toJavaObject(RequestMessage.class)); event.setResponseMessage(obj.getJSONObject(responseMessage).toJavaObject(ResponseMessage.class)); event.setTimestamp(obj.getLong(timestamp)); return event; } Override public boolean isEndOfStream(LogEvent nextElement) { return false; } Override public TypeInformationLogEvent getProducedType() { return TypeInformation.of(LogEvent.class); } }isEndOfStream返回false表示这是无界流Kafka 数据源源不断。getProducedType返回类型信息Flink 的序列化器需要用它来推断运行时类型。这里有个高频翻车点如果LogEvent内部有泛型字段TypeInformation.of会拿到泛型擦除后的类型导致序列化异常。解决方法是构造TypeInformation时带上泛型参数不过这个包里的类没有泛型嵌套用of没问题。3.3 计算链路中的 KeyBy 与 WindowLogEventApp里应该有keyBy和window的调用这是整个作业的计算核心。日志场景最常见的计算是按请求路径或接口维度统计请求量、延迟均值。窗口类型的选择直接影响结果语义我反推的典型逻辑如下stream .keyBy(event - event.getReqInfo().getPath()) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .aggregate(new CountAggregate()) .addSink(redisSink);keyBy按请求路径分组TumblingEventTimeWindows开一个 1 分钟的滚动窗口aggregate做增量聚合。滚动窗口适合每分钟统计一次滑动窗口适合做延迟敏感的指标比如每 10 秒更新一次最近 5 分钟的平均延迟。用事件时间窗口时水位线决定窗口什么时候触发而不是系统时钟这也是为什么前面说水印参数那么重要。3.4 消费性能与并行度配置Kafka Source 的并行度取决于两个因素一是 Kafka 主题的分区数二是env.fromSource之后 Transform 算子的并行度。Kafka Source 的每个并行子任务会消费至少一个分区分区数少于并行度时部分子任务会空闲。常见的设置方式是env.setParallelism(4)或者提交作业时通过-p 4参数指定。如果 Kafka 分区数是 12并行度设 4每个子任务消费 3 个分区。日志场景下单分区消费能力通常在每秒几千到几万条具体取决于消息大小和反序列化成本。消息体积大时瓶颈往往在反序列化而不是网络 IO这时候可以合并小消息批量解析或者简化LogEventSchema里的 JSON 解析逻辑。4. 双 Sink 落地Redis 集群写入与 MySQL 批量入库的参数细节4.1 Redis 集群 Sink 的客户端选型class 清单里没有直接出现 Redis 相关类名但摘要明确写了数据要导入 Redis 集群。Flink 官方没有独立的 Redis Connector社区常用的是flink-connector-redis它底层依赖 Jedis。Redis 集群模式下有个关键区别不能像单机那样直接指定redis://host:port必须用JedisCluster的节点列表方式。我一般推荐的写法如下FlinkJedisPoolConfig jedisConfig new FlinkJedisPoolConfig.Builder() .setHost(10.0.0.21) .setPort(6379) .setPassword(your-redis-password) .setDatabase(0) .setTimeout(3000) .build(); DataStreamString resultStream stream.map(LogEvent::toMetricString); resultStream.addSink(new RedisSink( jedisConfig, new RedisMapperString() { Override public RedisCommandDescription getCommandDescription() { return new RedisCommandDescription(RedisCommand.SET); } Override public String getKeyFromData(String data) { JSONObject obj JSON.parseObject(data); return obj.getString(metricKey); } Override public String getValueFromData(String data) { JSONObject obj JSON.parseObject(data); return obj.getString(metricValue); } } ));FlinkJedisPoolConfig是连接池配置setTimeout设置毫秒级连接超时。集群模式要换用FlinkJedisClusterConfig传入多个节点地址客户端会通过MOVED重定向找到正确的槽位不需要业务侧感知槽分配。.setDatabase(0)只在单机或哨兵模式下有意义集群模式不支持选择 database写了也会被忽略。4.2 集群 vs 单机的写入差异Redis 集群的生产者消费者行为跟单机有几个重要差异。第一KEYS命令在集群里不可用因为KEYS会扫描全部节点代价极高Flink Sink 里不能用这个命令做数据校验。第二SET操作按 key 的 CRC16 哈希值路由到对应槽位所以写入是自动分布的不需要手动指定节点。第三集群模式下MGET只能命中同一槽位的 key跨槽位的批量读取会报错如果需要批量读需要把相关的 key 设计成带同一个哈希标签比如{user:123}:profile和{user:123}:orders。4.3 MySQL JDBC Sink 的批量写入配置MySQL 写入用的是JdbcOutputFormat或者JdbcSink核心参数是批量大小和事务控制。Flink 的 JDBC Sink 不是来一条写一条那样吞吐太差。我的惯用配置是把 batch size 压在 1000 到 5000 之间连接参数也要跟上JdbcExecutionOptions execOptions JdbcExecutionOptions.builder() .withBatchSize(2000) .withBatchIntervalMs(200) .build(); JdbcConnectionOptions connOptions JdbcConnectionOptions.builder() .withUrl(jdbc:mysql://10.0.0.31:3306/log_analysis) .withDriverName(com.mysql.cj.jdbc.Driver) .withUsername(flink_writer) .withPassword(your-password) .build(); stream.addSink(JdbcSink.sink( INSERT INTO request_stats (path, cnt, avg_latency, window_end) VALUES (?, ?, ?, ?) ON DUPLICATE KEY UPDATE cnt VALUES(cnt), avg_latency VALUES(avg_latency), (ps, event) - { ps.setString(1, event.getPath()); ps.setLong(2, event.getCount()); ps.setDouble(3, event.getAvgLatency()); ps.setLong(4, event.getWindowEnd()); }, execOptions, connOptions ));withBatchSize(2000)表示攒满 2000 条执行一次批量写入withBatchIntervalMs(200)是一个兜底机制——即使数据量没到 2000超过 200 毫秒也会把当前批次刷出去。ON DUPLICATE KEY UPDATE处理的是窗口结果重复写入的场景Flink 的精确一次不能保证写入幂等需要 MySQL 端配合唯一键去重。注意ps.setTimestamp和ps.setLong的选择如果 SQL 字段是datetime类型传入毫秒时间戳需要先转成java.sql.Timestamp否则会丢时间精度或者直接报错。4.4 双写一致性怎么平衡Redis 和 MySQL 双写时两边数据并不是严格一致的。Redis 里的数据是实时的中间结果给看板和大屏用MySQL 里是持久化的统计结果给报表和离线分析用。Flink 作业对这两个 Sink 是独立写入的Redis 写入失败不会回滚 MySQL 的写入。生产中我一般接受这个不一致因为两个存储的用途不同但要注意把可重试的错误类型区分开Redis 连接超时可以重试MySQL 主键冲突不能重试否则无限重试会把消息堆积在算子后置队列里。5. 避坑指南Kafka 到 Redis/MySQL 链路的四个真实翻车点5.1 反序列化失败导致作业无限重启现象作业运行几分钟后进入反复重启循环Kafka 消费位点不前进Checkpoint 一直失败。原因Kafka 的某个分区里混入了一条非 JSON 格式的消息LogEventSchema.deserialize抛异常Flink 默认会把这个异常当成作业失败处理导致整个作业重启。重启后消费同一批数据再次失败死循环。解决在deserialize外层加 try-catch解析失败的消息先写出到死信队列或者直接丢弃但一定要记录日志和统计计数。另一个方案是在 Kafka 生产端加消息格式校验但消费端的兜底更稳妥。我用过的最稳妥方式是 catch 后把原始字节写入侧输出流后续对账时能查到那条脏数据长什么样。5.2 Redis 集群密码配置被忽略现象单机 Redis 写入正常换成 Redis 集群后所有写入报NOAUTH Authentication required。原因集群模式的鉴权是在每个节点上独立校验的而连接池配置里虽然设置了密码但JedisCluster初始化时没有把密码传给每个连接。很多旧版本的flink-connector-redis对集群模式的支持本来就不完善密码透传存在 bug。解决确认使用的连接器版本对JedisCluster的密码支持是完整的或者直接升级到 Jedis 4.x 的JedisCluster构造函数。如果升级不方便可以在 Redis 集群前面加一层哨兵或代理改为哨兵模式连接哨兵模式下密码传递比集群模式稳定得多。5.3 MySQL Sink 的写入字段类型不匹配现象作业整体运行正常但某个窗口结果写入 MySQL 时偶发报错错误信息是Data truncation: Out of range value或者Incorrect datetime value。原因Flink 的map或aggregate算子里输出的数值溢出比如avg_latency在某个极端窗口超过了 MySQLDOUBLE的精度范围或者时间戳字段被当成字符串拼进了 SQL 导致格式不对。解决在 Sink 的 SQL 绑定里显式做一次类型转换时间戳字段用new Timestamp(event.getWindowEnd())包装数值字段先做范围校验。另外检查 MySQL 表结构DOUBLE改成DECIMAL(10, 4)时间字段统一用BIGINT存毫秒时间戳能省掉一半的格式问题。5.4 作业重启后数据重复写入现象作业发生故障重启后MySQL 里的统计结果出现重复记录同一个窗口的数据写了两遍。原因Flink 的 Checkpoint 恢复只能保证算子状态恢复不能保证外部系统写入的幂等。Source 消费位点回退了窗口会重新计算然后 Sink 把结果再次写入 MySQL如果没有唯一键约束就会出现重复。解决在 MySQL 目标表上建立复合唯一索引字段就是窗口结果的自然键比如(path, window_end)。SQL 里用INSERT ... ON DUPLICATE KEY UPDATE保证重复写入时走更新而不是新增。Redis 侧用SETEX覆盖写入天然幂等不用额外处理。从那以后我每次设计双写链路都先问一句这个存储有没有自然键没有就先建键再做 Sink。6. 验证与进阶用数据对比确认整条链路的质量资源拿到手验证它是关键。class 包不能直接跑但可以把核心逻辑搭出一个最小可复现实验验证链路质量是否达标。我习惯先跑一个“冒烟验证”写一个生成器向 Kafka 发送 10 万条模拟日志时间戳按顺序生成故意让其中 5% 乱序 5 到 15 秒然后启动改造后的作业看三个指标。第一个指标是窗口触发延迟。事件时间窗口的触发时间等于windowEnd watermark如果水印允许乱序 10 秒窗口触发会比处理时间晚约 10 秒。拿 Redis 里结果的时间戳和当前时间对比偏差在可接受范围内说明水印配置基本合理。第二个指标是数据完整率。比对 Kafka 消费总数和 MySQL 写入总数用SELECT COUNT(*)对账偏差为零基本可以确认没有丢数据。第三个指标是写入速率。观察 Redis 和 MySQL 的写入 QPS如果 MySQL 明显低于 Kafka 消费速率说明batchSize太小或者连接池过小。进阶用法里最实用的是把 Sink 从异步改成多路并行写。Redis Sink 和 MySQL Sink 各有自己的连接池addSink之后 Flink 会自动做反压控制但如果想提升吞吐可以把两个 Sink 拆成两个DataStream分支分别设置并行度Redis 写路径并行度可以高一些MySQL 写路径受数据库连接数限制并行度一般不要超过数据库配置的连接数上限。我之前在这个场景踩过最深的坑是不管 Redis 还是 MySQL都用默认并行度跑结果 MySQL 被连接数打满Redis 反而闲着。从那以后我做双 Sink 作业第一件事就是分别看两个下游的写入能力再反推各自的并行度。验证无误之后这份 class 包就变成我手头最稳定的参照链路每次写新的 Flink 日志作业都会拿它对照一遍。希望帮到你。本文还有配套的精品资源点击获取
返回列表