☰
Flink读取Kafka数据:双写Redis与MySQL的实时链路实践
2026/10/7 19:01:54 网站建设 项目流程

简介:面向大数据流计算开发者和实时数仓初学者,这是一个围绕 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、LogEventWaterMarkExtractor
  • ReqInfo、RequestMessage、ResponseMessage
  • PropertiesConfiguration

这组命名说明,核心事件模型是一个日志事件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的调用点,说明作业用的是事件时间。日志类数据天然带客户端时间戳,用事件时间计算出的延迟指标才有业务意义,否则跑出来的数字全是服务器处理时刻,参考价值大打折扣。

这里我还原一下典型的水印提取写法,这个也是我在真实项目里常用的模式:

DataStream<LogEvent> withWatermarks = stream .assignTimestampsAndWatermarks( new LogEventWaterMarkExtractor() ); public class LogEventWaterMarkExtractor extends BoundedOutOfOrdernessTimestampExtractor<LogEvent> { 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"); KafkaSource<LogEvent> source = KafkaSource.<LogEvent>builder() .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(); DataStream<LogEvent> 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 DeserializationSchema<LogEvent> { @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 TypeInformation<LogEvent> getProducedType() { return TypeInformation.of(LogEvent.class); } }

isEndOfStream返回false表示这是无界流,Kafka 数据源源不断。getProducedType返回类型信息,Flink 的序列化器需要用它来推断运行时类型。这里有个高频翻车点:如果LogEvent内部有泛型字段,TypeInformation.of会拿到泛型擦除后的类型,导致序列化异常。解决方法是构造TypeInformation时带上泛型参数,不过这个包里的类没有泛型嵌套,用of没问题。

3.3 计算链路中的 KeyBy 与 Window

LogEventApp里应该有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(); DataStream<String> resultStream = stream.map(LogEvent::toMetricString); resultStream.addSink(new RedisSink<>( jedisConfig, new RedisMapper<String>() { @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 日志作业都会拿它对照一遍。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询