做了这么多年数据开发,我越来越觉得,大数据仓库最值得琢磨的不是某个框架有多酷,而是怎么把离线批处理和实时流计算真正“熔”到一块儿。Flink 和 Hive 的集成,就是目前落地批流一体最顺手的路子之一:Hive 继续当数仓底座,负责元数据、分区和最终存储,Flink 负责把实时数据写进来、把离线计算跑得更活。
这套组合几乎涵盖了一个数据平台最常提的几个热词:批流一体、实时数仓、数据湖底座。但网上讲 Flink 的文章很多,讲 Hive 的文章也很多,专门把两者“接缝”处讲透的却不多。我见过太多同学,Flink 单独跑得很溜,Hive 也都熟练,一到集成就卡在类加载、元数据隔离、方言不兼容上。这篇文章就围绕这些实际内容展开,适合正在搭数仓但被两套体系折磨的工程师,想搞明白 Flink SQL 与 Hive 怎么配合的转型选手,以及被小文件、JDBC 连接器异常这类问题反复折腾过的人。
我尽量不堆理论,只讲能直接抄作业的思路和步骤。版本匹配、Catalog 和方言、流式写入 Hive、小文件治理、MySQL 同步 ClickHouse、Spring Boot 集成 Flink,都会过一遍。这些都是真实项目里高频出现的点,读完之后你应该能少踩几个坑。
1. 批流一体到底解决什么问题
1.1 批与流的“两张皮”怎么来的
在 Flink 和 Hive 集成成熟之前,大多数数仓是这么跑的:业务日志先进 Kafka,实时链路用 Flink 做清洗,然后写进 Redis、ES、ClickHouse 供大屏和在线服务查询;离线链路等凌晨由 Hive 或 Spark 再消费一遍 Kafka,或者从业务库直接抽数,落成 Hive 分层表。白天看实时看板,晚上看离线报表,看起来互不干扰,实际上暗坑无数。
同一个用户指标,实时表里按最近 2 小时的行为定义,离线表里按自然日去重定义,数值对不上,业务来回找数仓。两套链路用的 schema 各自维护,字段名改了,实时作业和离线任务各崩一次,事故都能双份。更糟的是消息队列重放的时候,实时侧可能已经提前把结果写进去了,离线侧因为还没到凌晨,根本不知道这条数据后来变了。
批流一体的最初动力,说白了就是不想再维护两套口径。它也不是要把 Lambda 架构推翻重来,而是把最容易扯皮的表结构、字段口径、存储底座统一起来,让局部分区级仍用批处理、业务需要实时时用流计算。Flink 和 Hive 之间的集成,正好提供了这种统一的外壳。
1.2 Flink 和 Hive 之间那条“缝”在哪
Flink 与 Hive 集成的本质,是让 Flink 看见 Hive 的元数据、用 Hive 的方言写 SQL、把数据直接落到 Hive 的分区里,同时还能像读流一样感知 Hive 新产生的分区。拆开看有三条关键线索:
- 元数据线:HiveCatalog 把 Hive 的数据库、表、分区、视图映射成 Flink 的 Catalog 对象,Flink 不用再造一套 schema。
- 方言线:Flink SQL 默认方言是 Flink,但可以切成 Hive 方言,让那些在 Hive 里写了几年的 SQL 直接跑在 Flink 上。
- 数据线:Flink 可以把 Kafka 的实时数据直接插进 Hive 表,也可以从 Hive 表新增分区触发流式计算,实现真正意义上的“批流复用同一张表”。
有朋友问我:直接用 Hudi、Iceberg 不是更洋气?确实洋气,但在大部分公司,Hive 仍然是被安全、审计、BI、报表工具都认的那套标准。Hudi 和 Iceberg 有它们自己的优势,但也有迁移成本和学习成本。Flink 先跟 Hive 打通,相当于把新引擎的能力挂到了老体系的接口上,前期见效最快,后面真想换湖格式,Hive 表也能平滑过渡。所以如果你还没有非得用湖格式的理由,先从 Hive 集成开始是比较稳妥的选择。
2. 集成前必须先做好的环境与版本匹配
2.1 版本选型的现实经验
接触过 Flink 和 Hive 集成的人都知道,最大的坑通常不是用法,而是版本。不同 Flink 版本对 Hive 版本的支持范围不一样,官方兼容矩阵一定要查,不要看到网上某个配置能跑就照抄。
我自己常用的组合是 Flink 1.17.x 配 Hive 3.1.3,Spark 侧如果要共用同一套 Hive Metastore,也能兼容。选择这个组合的原因很简单:Hive 3.1.2/3.1.3 是目前生产环境最普及的稳定版本,Flink 官方发布的flink-sql-connector-hive-3.1.3_2.12对应 jar 也是现成的,省去手动拼依赖的麻烦。如果你们平台还在用 Hive 2.x,也请走 Flink 官方标明的对应连接器,不要混装。
除了 Flink 和 Hive 的版本,还要注意 Hadoop 和相关 jar 的版本。Flink 运行时默认带一套 Hadoop 兼容层,但最好还是把集群的hadoop-client依赖放进去,否则访问 HDFS 上的高层文件时容易出现权限或协议不匹配。如果 Hive Metastore 元数据库用的是 MySQL,通常还需要补一个 MySQL JDBC 驱动,这个点经常被漏掉,导致 Catalog 连不上。
我的建议是先在你的测试机上搭一个单节点 Hive Metastore,把目标版本确认好、把连接跑通,再推给集群。直接在生产环境试版本,出了错根本分不清是网络问题还是依赖问题。
2.2 把 Hive 的依赖送进 Flink
官方推荐方式是把对应的flink-sql-connector-hivejar 放进 Flink 的lib目录。文件路径大概长这样:flink-sql-connector-hive-3.1.3_2.12-1.17.2.jar。放进lib后,需要重启 Flink 集群或 SQL 客户端才能加载。
另一个常见做法是把 Hive 的安装目录也加入依赖,但我们要注意别以类冲突的方式来加。更稳的方式是通过HADOOP_CLASSPATH注入 Hadoop 和 Hive 相关类,Flink 脚本会读取这个环境变量。在提交任务的节点上,把下面的配置写进~/.bashrc或提交脚本:
export HADOOP_HOME=/usr/local/hadoop export HADOOP_CLASSPATH=$(hadoop classpath) export HIVE_HOME=/usr/local/hive export HIVE_CONF_DIR=/usr/local/hive/conf然后把hive-site.xml复制到$FLINK_HOME/conf目录。原因是 Flink 的 HiveCatalog 在初始化时,会优先从hive-conf-dir参数指定的目录读取配置,其次从FLINK_HOME/conf下找hive-site.xml。你若没有显式指定hive-conf-dir,不把配置文件放对位置,它就找不到 Metastore 地址,报一堆连接失败。
2.3 启动 SQL 客户端验证连通
环境准备好了,先用 SQL Client 做一个最小验证。启动bin/flink-sql,执行:
CREATE CATALOG hive_catalog WITH ( 'type' = 'hive', 'hive-conf-dir' = '/usr/local/hive/conf', 'hive-version' = '3.1.3' );然后执行:
USE CATALOG hive_catalog; SHOW DATABASES;能列出 Hive 里已有的库,说明 Metastore 连通了。再执行一条SHOW TABLES和一个简单的 Hive 表SELECT COUNT(*),如果 Hive 目录下有 ORC 或 Parquet 数据,读得出来,你的集成环境基本就算通了。这套验证不要省,后面所有问题排查,都比在集群故障时再发现基础依赖错误轻松得多。
3. 三个核心集成点:Catalog、方言、流式写入
3.1 HiveCatalog:让两套引擎只认一份元数据
很多人不知道为什么 Flink 需要专门搞一个 HiveCatalog。原因很简单:如果 Flink 自己维护一套表,Hive 自己维护一套表,两边根本不认识,pipeline 一重启就找不到表了。用 HiveCatalog 之后,Flink 里的hive_catalog.default.user_operation指的就是 Hive 里的同一个表,你在 Hive 里ALTER TABLE加一列,Flink 侧几乎无需改动就能感知到。
这带来一个很实际的好处:身份权限和安全策略在 Hive 侧已经定了,Flink 通过 HiveCatalog 访问时,只要提交作业的用户有对应权限,就能直接沿用。反过来说,作业提交用户没有某张表权限,Flink 也会报错,不会绕过 Hive 的权限体系。
需要注意的是,Flink 的 HiveCatalog 并不只是“读” Hive 的元数据;用 Flink SQL 在 HiveCatalog 下建表,如果开了 Hive 方言,它会直接把表建到 Hive 侧。这样既能在 Flink 里写任务,又能在 Hive 里用老 SQL 查,属于最标准不过的批流一体底座。
3.2 Hive 方言:让老 SQL 直接跑起来
Flink 默认的方言还是 Flink SQL,语法跟 Hive 有细微差别。为了适配 Hive 的 DDL、函数和LATERAL VIEW之类的语法,Flink 提供了 Hive 方言开关:
SET table.sql-dialect=hive;这个开关要放在每一个 session 一开始,尤其要注意,它会影响后续执行的 DDL 和查询。切换方言后,你可以直接使用CREATE TABLE、ROW FORMAT、STORED AS、TBLPROPERTIES这类 Hive 风格语法,写出来的表也会真正落到 Hive。
最常见的坑是有人把方言设置得不彻底:建表时用 Hive 方言,查询时又切回 Flink 方言。建议在一个作业内固定一种方言,除非你是故意想在 Flink 方言下用 Hive 的底表做实时查询。后者并不冲突,但你要心里有数:读已有的 Hive 表,不需要切方言;用 Hive 风格去建新表或执行 Hive 专属语法,才需要切。
3.3 流式写入 Hive:分区提交不是可有可无
把 Kafka 数据实时写进 Hive 分区表,听起来就是 insert 一张表,实际上涉及三个环节:分区路径怎么落到 HDFS、分区什么时候算写完、写完怎么通知下游。
Flink 写 Hive 表时,分区提交策略是最重要的参数。配置大致如下:
CREATE TABLE hive_events ( event_id BIGINT, event_name STRING, event_ts TIMESTAMP(3), WATERMARK FOR event_ts AS event_ts - INTERVAL '1' MINUTE ) PARTITIONED BY (dt STRING, hr STRING) STORED AS PARQUET TBLPROPERTIES ( 'sink.partition-commit.trigger' = 'partition-time', 'sink.partition-commit.delay' = '1 h', 'sink.partition-commit.policy.kind' = 'metastore,success-file' );partition-time表示用事件时间来判断分区是否到达提交临界点,delay则缓冲一部分迟到的数据;metastore会把分区信息写进 Hive Metastore,success-file则会在分区目录里生成一个_SUCCESS文件,Hive 或下游任务看到这个文件就知道分区可用了。如果你用process-time,则不需要管事件时间,但数据迟到时容易把数据写到上一个分区,后续查数时就对不齐。
还有一点很多教程不强调:流式写 Hive 时,分区字段不要用字符串硬拼,时间字段尽量在源表里用TIMESTAMP类型,并且给上WATERMARK。这样分区提交才有一个可靠依据。否则一切靠“作业进程看到时钟到了”来提交,恢复或重放时问题会特别多。
4. 实战:Kafka 实时入仓 + 小文件治理
4.1 从 Kafka 到 Hive 的完整 SQL 链路
先给 Hive 建一张目标分区表,字段按业务明细定义,存储格式用 Parquet:
CREATE TABLE app_event ( event_id BIGINT, event_name STRING, event_time TIMESTAMP ) PARTITIONED BY (dt STRING, hr STRING) STORED AS PARQUET;然后在 Flink SQL 里定义 Kafka 源表。注意事件时间字段要声明水位线,否则后面的流式分区提交没有依据:
CREATE TABLE kafka_event ( event_id BIGINT, event_name STRING, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '1' MINUTE ) WITH ( 'connector' = 'kafka', 'topic' = 'app_event', 'properties.bootstrap.servers' = 'kafka-1:9092,kafka-2:9092', 'properties.group.id' = 'flink_dw_group', 'scan.startup.mode' = 'latest-offset', 'format' = 'json', 'json.timestamp-format.standard' = 'ISO-8601' );接下来注册 HiveCatalog,并直接向 Hive 表写入:
CREATE CATALOG hive_catalog WITH ( 'type' = 'hive', 'hive-conf-dir' = '/usr/local/hive/conf' ); USE CATALOG hive_catalog; INSERT INTO default.app_event SELECT event_id, event_name, event_time, DATE_FORMAT(event_time, 'yyyy-MM-dd') AS dt, DATE_FORMAT(event_time, 'HH') AS hr FROM default.kafka_event;这里分区字段由 Flink SQL 自己在写入时计算,比在 Hive 侧触发MSCK REPAIR再去发现目录要干净得多。如果哪天你已经在 Kafka 里存了整天的数据,想回刷历史分区,把scan.startup.mode改成earliest-offset再跑一遍即可。
4.2 流式提交参数与时区陷阱
流式入仓最常出问题的不是 SQL 本身,而是时区。假设 Kafka 里的event_time是 UTC 时间,而业务希望按上海时区分区,那 SQL 里直接DATE_FORMAT很容易把早上 8 点前的事件归到前一天,早上报表先炸。
处理办法有两种。第一种是数据源尽量规范,上游在写入 Kafka 前统一成东八区时间;第二种是在 Flink 端显式转换,比如:
SELECT event_id, event_name, CONVERT_TZ(event_time, 'UTC', 'Asia/Shanghai') AS event_time_local, DATE_FORMAT(CONVERT_TZ(event_time, 'UTC', 'Asia/Shanghai'), 'yyyy-MM-dd') AS dt, DATE_FORMAT(CONVERT_TZ(event_time, 'UTC', 'Asia/Shanghai'), 'HH') AS hr FROM kafka_event;另外,Flink 1.15 以后有一个table.local-time-zone配置,默认跟随系统时区。如果集群和业务时区不一致,建议在flink-conf.yaml里统一改成Asia/Shanghai,否则TIMESTAMP_LTZ类型在查询和写入时很容易出现偏移。
4.3 小文件治理:别等 HDFS 告警才想起来
流式写入分区表的场景下,小文件几乎是必然。原因很简单:Flink 默认并行度可能很高,每个 subtask 一次提交只写一个 part 文件;分区提交频率越高,小文件越多。Hive 底表一旦太多小文件,查询时 NameNode 压力大,扫描任务打开文件数量爆炸。
治理三板斧。
一是在 Flink 侧降低文件碎片:调整并行度、增加sink.partition-commit.delay、设置文件滚动参数。对于文件系统连接器,可以控制单文件滚动大小,类似:
'sink.rolling-policy.file-size' = '128MB', 'sink.rolling-policy.rollover-interval' = '10 min'如果用了 Hive 连接器,部分连接器参数也可以直接以TBLPROPERTIES方式传给下层连接器,具体以你使用的 Flink 版本文档为准。
二是用批式重写去压缩:等一个分区不再写入后,执行INSERT OVERWRITE,把该分区重写成全量一个或少数几个文件。类似:
SET hive.exec.dynamic.partition.mode=nonstrict; INSERT OVERWRITE TABLE app_event PARTITION (dt, hr) SELECT event_id, event_name, event_time, dt, hr FROM app_event WHERE dt = '2025-01-01' DISTRIBUTE BY dt, hr;DISTRIBUTE BY非常关键,它会把落在同一个分区的数据尽量分给同一个 reducer,避免写完还是几十个小文件。
三是给 Hive 设合并参数,针对小文件多的表,在跑批任务里加:
SET hive.merge.mapfiles=true; SET hive.merge.mapredfiles=true; SET hive.merge.size.per.task=268435456; SET hive.merge.smallfiles.avgsize=134217728;还可以对 ORC 或 Parquet 表执行ALTER TABLE app_event CONCATENATE;。这类操作会把表目录下的小文件合到一起,不会改写数据内容。注意 textfile 格式不能用CONCATENATE,需要先转成列式格式再合并。
5. 热搜里的真实场景:JDBC 异常、MySQL 到 ClickHouse、Spring Boot 集成
5.1 Flink JDBC 连接器为什么会抛异常
为什么搜“flink的jdbc连接器异常”的人这么多?太真实了,因为十个 Flink JDBC 作业九个是栽在连接器上。拿常见的ClassNotFoundException: com.mysql.cj.jdbc.Driver来说,八成是 JDBC 驱动 jar 没打进 lib,或者 fat jar 里把多个 MySQL 驱动版本混了。检查顺序按三个方向来:
- 依赖有没有:把
mysql-connector-j或mysql-connector-java的实际 jar 放进 Flinklib或任务 jar,并保证多个 jar 里不冲突。 - 参数对不对:MySQL 8 的驱动类名是
com.mysql.cj.jdbc.Driver,URL 里要带时区参数,比如jdbc:mysql://rm-xxx.mysql.rds.aliyuncs.com:3306/db?useSSL=false&serverTimezone=Asia/Shanghai。 - 连接要不要复用:Flink JDBC sink 默认会管理连接池,但如果你在自定义 sink 里频繁
getConnection,连接池被耗尽就报Too many connections。
还有一种很像连接异常的错误:SQL 写好了,但目标表字段类型对不上。比如 MySQL 的datetime传到 Flink 里变成了TIMESTAMP_LTZ,目标端期望的是String,连接器抛错。建议在同步前先SHOW CREATE TABLE看清类型,再用CAST显式转换。
5.2 用 Flink 实现 MySQL 同步到 ClickHouse
这个场景如今比同步到 ES 还常见。核心思路是用 Flink CDC 抓 MySQL 的 binlog,结构化之后写进 ClickHouse。ClickHouse 不是事务型数据库,默认没有行的 update/delete 语义,所以行级同步需要借助ReplacingMergeTree或AggregatingMergeTree去重版本。CDC 数据里一般带op字段,删除事件要谨慎,不要把删除直接当成新数据写进去。
一种常用写法是:
- 源表:
mysql_cdc连接器,scan.startup.mode=initial,读存量后再抓增量。 - 目标表:ClickHouse
ReplacingMergeTree引擎,表结构里加一个版本字段,直接映射 binlog 的ts_ms。 - Flink 端:JDBC sink 或 ClickHouse HTTP sink,写入时
insert into ck_table select ... from mysql_cdc_source。
同步链路里最容易踩的问题是压力写:ClickHouse 单次大批量 insert 性能很好,但 JDBC connector 默认缓冲参数要调。一般配sink.buffer-flush.max-rows=1000、sink.buffer-flush.interval=5s,并行度不要开太多,否则 ClickHouse 容易被“打死”。数据量更大时,建议走官方 HTTP 批量提交端口或分布式表写入,不要全部挤在 JDBC 上。
5.3 Spring Boot 集成 Flink 的正确姿势
Spring Boot 能集成 Flink,但做这种事情前先想清楚部署形态。如果只是本地调试,Spring Boot 里搞一个StreamExecutionEnvironment,把 SQL 跑一遍没问题,问题在于它启动的其实是一个内嵌的 Flink 客户端,作业的生命周期跟着 Spring Boot 进程走。进程一重启,所有作业都没了。所以线上不建议把 Flink 环境写在 Spring Boot 应用里,除非你只跑开发测试。
实际生产常用的是“Spring Boot 负责管理任务配置,远程 Flink 集群负责跑”的分离模式。Spring Boot 端把 SQL 模板、源表配置、目标表配置管好,需要提交时通过 Flink REST API 把 jar 提交到 JobManager,或者调用flink run的封装接口。这个模式好处是应用重启、发版都不影响作业,坏处是你要额外维护一套作业提交协议和状态管理。
如果一定要在 JVM 进程里跑一个小型 Flink 任务,至少要把类加载模式搞清楚。Spring Boot 默认的 fat jar 结构容易跟 Flink 的lib依赖冲突,提交任务时会出现各种NoClassDefFoundError或序列化器报错。我的处理办法是构建一个专门的flink-runner模块,只依赖flink-table-api-java和所需 connector,用maven-shade-plugin打包,把 Spring Boot 的类隔离在外。
5.4 Hive 窗口函数在数仓里的高频用法
热点词里还有一个“hive给每一行标号”,对应就是窗口函数。用户在数仓里最基本的诉求是给明细表加一个行号,比如取每台设备最近一条登录记录。Hive 的写法是:
SELECT device_id, login_time, ROW_NUMBER() OVER (PARTITION BY device_id ORDER BY login_time DESC) AS rn FROM login_log;外层再包一层WHERE rn = 1,就能得到每台设备的最新登录时间。这个需求看着简单,实际在离线数仓里就是“去重取最新”的标准答案。
Flink 集成 Hive 后,这套窗口函数基本上也能直接用。如果你切到 Hive 方言,甚至不用改语法;如果你用 Flink 方言,语法也接近,只是部分函数名略有差异。需要注意ROW_NUMBER在流计算里会产生无界窗口的状态开销,如果数据量极大且乱序严重,要给ORDER BY的字段加水位线控制,否则状态无限增长,作业迟早内存爆掉。离线批处理就没这个问题,这其实是批流一体最现实的一面:同一个 SQL,在两种模式下面临的另一套物理问题。
6. 实操中必须记住的坑和经验
6.1 常见问题速查表
这里整理一份我日常排查问题最常用的速查表,基本能覆盖 Flink 和 Hive 集成的大部分初期故障。
| 症状 | 可能原因 | 处理建议 |
|---|---|---|
| HiveCatalog 建不出来 | 没放 Hive 连接器 jar,或hive-site.xml不在 conf 目录 | 检查lib下是否有对应 connector jar,确认hive-conf-dir路径 |
| 连接 Metatore 报权限或超时 | Metastore 地址配置错误,或 MySQL 驱动缺失 | 核对hive-site.xml的thrift地址,补 MySQL JDBC 驱动 |
| 流式写 Hive 没有新分区 | 分区提交策略没配,或没开metastore提交 | 设置sink.partition-commit.policy.kind并检查作业日志 |
| 表里全是小文件 | 并行度过高,分区提交太频繁 | 调并行度,设置滚动策略,定期CONCATENATE或重写分区 |
| JDBC 同步连接不上 | 驱动缺失、时区参数不对 | 加 jar,URL 后面补serverTimezone,确认驱动类名 |
| ClickHouse 同步数据对不上 | 删除事件处理不对,缺去重键 | 用ReplacingMergeTree,把删除转成版本字段或标记字段 |
| Spring Boot 提交后作业丢失 | 作业生命周期绑定在 Spring Boot 进程上 | 改成远程提交模式,Spring Boot 只做配置管理 |
| Flink 和 Hive 字段类型不匹配 | 两边 schema 没有对齐 | 建表前先SHOW CREATE TABLE,用CAST做显式转换 |
| 分区时间少 8 小时 | 集群时区不是东八区 | flink-conf.yaml里指定table.local-time-zone为Asia/Shanghai |
6.2 我在生产环境踩明白的几个经验
第一,能少引依赖就少引依赖。Flink 集成 Hive 之后,classpath 已经够复杂了,再随便把几个连接器全塞进 lib,经常出现不同 jar 里的同一类冲突。我的习惯是只用官方 connector,Hadoop 和 Hive 的类尽量走HADOOP_CLASSPATH,不要手动拼一堆 jar。
第二,流式写 Hive 的任务,一启动不要直接开高并行度。很多人在测试环境用几行数据把 SQL 跑通,上生产直接给 32 并行,结果一个分区写了一百多个小文件,最后还得花更多时间治理。我一般先用低并行度跑一周,统计分区的平均文件大小,再决定要不要加并行。数据量真的上来了,优先加的也是并行度而不是分区提交频率。
第三,分区提交要能随时看到。我把_SUCCESS文件当作生产可用的信号以后,下游的任务都改成只扫描带成功标记的分区,这样至少能避免数据没写完就被离线任务读到。配合一个定时的分区大小巡检,每天凌晨检查前一天写了多少个文件、每个文件多大,超过阈值就自动触发一次INSERT OVERWRITE重写。这套流程不需要人肉盯,但能把小文件问题消灭在报表发布之前。
我记得第一次给业务方演示批流一体时,对方并不关心你用了什么技术,只问:昨天中午改的口径,现在实时报表和日报能不能一致?当时我用的方案就是 Kafka 实时写 Hive 分区,Hive 表同时被离线任务读取,Flink 用同一套 Catalog 和方言处理两条链路。那次演示的指标确实对上了,后来我总结出的道理特别朴素:批流一体不是给技术看的表演,就是让别人按同一张表反复查数,结果不打架。能做到这一点,集成就算成功。