ARTICLE DETAIL

资讯详情

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

Flink与Hive集成实战:批流一体落地与分区提交小文件治理

Flink与Hive集成实战:批流一体落地与分区提交小文件治理 1. 为什么需要把Flink和Hive放在一起批流一体的真实痛点先聊一个每天都在发生的场景。公司的数据链路往往是这个样子业务库的binlog被CDC工具抓出来进KafkaFlink消费后实时写入ClickHouse、Doris或者HBase供前端查询另一边ODS层的原始日志和业务数据用Hive SQL做T1清洗落到DWD、DWS层再产出报表。两条链路并行存在代码各写一套口径靠人肉对齐数据从实时到离线还要再做一次同步。维护成本高不说最怕的是“实时算出来的日活”和“离线跑出来的日活”对不上线上排查一整天最后发现是窗口中位数和状态保留策略的差异。所谓的“批流一体”就是不想再维护两套引擎、两套口径、两套存储。Flink本身在1.12之后已经把批流API统一了DataStream和Table API既能跑有界流也能跑无界流但这只是计算层的事情。真正让批流一体落地的关键是存储和元数据层也能打通。而大多数公司数仓的底座就是Hive所以Flink和Hive的集成就成了绕不开的一步让Flink能直接读Hive表、写Hive表、复用Hive的MetaStore和函数让实时任务产出的数据能无缝进入离线数仓体系也被离线团队直接消费。我见过不少团队在初期只是把Flink当“数据搬运工”从Kafka到HDFS或者从MySQL到Hive等业务提出“我要看实时昨天的32个指标口径和离线报表完全一致”时才开始认真研究Flink和Hive的深度集成。这篇文章就把我实际踩过的坑和沉淀下来的可用方案完整写一遍覆盖Hive Catalog、流式写入、分区提交、小文件治理、依赖冲突和性能调优适合正在搭实时数仓或者想把离线实时两套链路合并的工程师。2. 集成方案选型Hive Catalog、Hive Streaming与Hive Dialect2.1 Hive Catalog到底是什么Flink和Hive集成核心是通过HiveCatalog来绑定Hive的MetaStore。你可以把MetaStore理解成数仓的“户籍系统”里面记录了有哪些库、哪些表、字段是什么、分区有哪些、数据文件在HDFS的哪个路径。Flink里创建HiveCatalog之后就能把Hive表当作Flink的表来用Flink SQL里可以直接写CREATE TABLE xxx (...) WITH (connectorhive ...)也可以直接用USE CATALOG myhive;切换到Hive Catalog下操作。这个设计省掉了大量“填WITH参数”的体力活。以前用Flink读写Hive每个人都要手动指定path、format、partition等一堆参数Catalog机制则直接从MetaStore拉取Schema底层文件格式、字段类型、分区路径全部自动映射。更重要的是Catalog保证了Flink看到的表结构和Hive完全一致不会再出现“离线表字段顺序变了实时写进去错位”的经典翻车。Flink的HiveCatalog是模块化设计的它不绑定特定Flink版本之外的东西但要正确工作需要引入对应Hive版本的flink-sql-connector-hive-x.x依赖。版本兼容矩阵我会在第三章写清楚这里先说结论选版本时以Flink官方文档的兼容表为准不要想当然地以为Hive 3.1.2和Flink 1.17必然兼容实际还要看JDK和Hadoop版本。2.2 Flink流式写入Hive的两种模式Flink实时写Hive两种典型做法一种是用StreamingFileSink其实在较新版本推荐FileSinkStreamingFileSink已被标记废弃配合HiveBulkWriter另一种是用Flink SQL的INSERT INTO配合hive连接器和streaming相关参数。第一种更适合DataStream API场景你在StreamExecutionEnvironment里读Kafka的DataStreamRow然后通过HiveBulkWriter或FileSink写入Hive分区目录。这种方式的优点是灵活能直接控制写文件的格式、滚动策略和分区目录结构但需要自己处理提交逻辑。第二种是多数项目会选的方式。Flink SQL里直接写INSERT INTO hive_catalog.db.ods_table SELECT ... FROM kafka_source然后通过表参数控制写Hive的行为比如设置streaming-source.enable、sink.partition-commit.trigger等。Flink会把每条数据按分区字段路由到对应的Hive分区目录写入的文件是PartFile需要等到checkpoint完成并且满足分区提交条件时才把_SUCCESS标记文件和分区信息注册到MetaStore。这一点非常关键Hive表在Flink写入过程中如果不做分区提交Hive侧始终看不到数据。网上搜“flink的jdbc连接器异常”很多案例就是在测试Flink写Hive时把jdbc连接器和hive连接器搞混。Flink写Hive用的是connector hive不是jdbc。JDBC连接器是给MySQL/PG这类支持JDBC协议的数据库用的如果遇到Could not find any factory for connector jdbc这类报错多半是缺了flink-connector-jdbc的依赖而不是Hive的问题。2.3 什么时候用Hive DialectFlink 1.13开始引入了Hive Dialect简单说就是让Flink SQL能够解析Hive的语法。比如Hive里有LATERAL VIEW、TRANSFORM、CLUSTER BY这类Flink原生SQL不支持的语法用Hive Dialect就能直接跑。但这个特性很容易被误解。我遇到过有人为了让Flink能写Hive表就把所有SQL都切换到Hive Dialect结果在流式任务里跑INSERT OVERWRITE语义完全不对。Hive Dialect设计的主要目的是离线批式场景下的SQL兼容不是给流式任务用的。流式任务请继续使用Flink原生Dialect只在执行离线批查询、批量修正数据、或者必须用Hive UDF的时候切到Hive Dialect。具体切换方式是在Flink SQL客户端执行SET table.sql-dialect hive;注意Dialect是会话级别的设置之后默认的CREATE TABLE和INSERT都会走Hive语法解析。所以在同一段脚本里混用原生SQL和Hive SQL时需要来回切换别嫌麻烦。2.4 选型建议如果只是离线批量读Hive表然后算点东西用Flink的批模式加HiveCatalog就够了连流式特性都不用开。如果要实时从Kafka消费数据写入Hive分区表走Flink SQL的CREATE TABLE ... WITH (connectorhive) 分区提交机制这是目前最稳的经典路线。如果要流读Hive表比如Hive表有新数据进来Flink能感知并当流来读就得开启streaming-source.enable并且注意Hive表必须是分区表且分区目录要有固定的命名模式。这里对Hive表的“目录结构规范”要求很高如果Hive表是外部表且分区路径乱建流读是玩不起来的。如果想把Hive的UDF直接在Flink SQL里用可以在Catalog里注册function保证Flink作业提交时能把Hive的auxlib或自定义UDF jar带上。这个实操里容易踩坑后面第五节专门讲。3. 环境准备与核心配置3.1 版本兼容矩阵先列我实测比较稳的组合FlinkHiveHadoop说明1.15.x2.3.92.7.5老集群稳定但Hive功能受限1.16.x2.3.9 / 3.1.33.1.0我目前生产环境的主力1.17.x3.1.33.1.0新特性多注意依赖冲突1.18.x3.1.33.1.0对JDK11更友好不推荐Hive 4.x和Flink集成因为Flink官方连接器长期只维护到3.1.xHive 4.x的MetaStore协议有些变化要用得自己打补丁运维成本高。另外Flink官方发布了两个关键的构件flink-sql-connector-hive-hive.version这是个“伪shaded”的jar里面包含了执行Hive方言和HiveCatalog所需的部分依赖另一个是flink-connector-hive它更轻量适合和Flink自带的Hadoop依赖配合使用。如果你用flink-connector-hive还得显式引入hive-exec、hive-metastore等Hive依赖麻烦不少。所以我个人建议直接用flink-sql-connector-hive省心。3.2 依赖打包与集群部署如果是用flink run提交作业直接把连接器jar放到Flink的lib目录下最省事。但我更推荐在项目里用Maven打包时引入依赖这样能控制版本避免和Flink内置依赖冲突。典型pom.xml片段dependency groupIdorg.apache.flink/groupId artifactIdflink-sql-connector-hive-3.1.3_2.12/artifactId version1.17.2/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-table-api-java-bridge_2.12/artifactId version1.17.2/version /dependency特别提醒flink-sql-connector-hive内部已经shaded了部分hadoop和hive依赖如果你在项目里又显式引入了hive-exec、hadoop-client等极易出现NoSuchMethodError或者ClassCastException。我踩过一次Flink作业一提交就报IncompatibleClassChangeError折腾了很久才发现是自己把hadoop-common的版本强依赖到了3.2和Flink内置的3.1.0冲突。解决方式很简单把项目里多余的hadoop依赖全部排除只留连接器。flink-sql-connector-hive和HiveCatalog的加载还依赖一个不显眼但必须的东西Hive的hive-site.xml。很多教程没说清楚没有这个文件Flink启动时只会在classpath里找默认配置找不到hive.metastore.uris就连不上MetaStore。所以必须保证作业运行环境中能找到hive-site.xml要么把它打进jar包要么在提交脚本里用-C把配置文件所在目录加到classpath。我在YARN上是用ship-files参数带过去的flink run -t yarn-per-job \ -Dyarn.ship-files/etc/hive/conf/hive-site.xml \ -c com.example.Job myjob.jar3.3 核心参数配置HiveCatalog构造时有两类核心参数一类是Hive本身的另一类是Flink的。Hive侧必配的参数是hive.metastore.uris格式是thrift://host:9083。注意老集群MetaStore可能用的不是默认端口或者配置了高可用这时候hive-site.xml里会有多个urisFlink会自动读取。强烈建议不要在代码里硬编码MetaStore地址最好只依赖hive-site.xml。Flink侧参数主要在创建HiveCatalog时传入比如HiveCatalog catalog HiveCatalog.create( new HiveCatalog.HiveCatalogBuilder() .setHiveConfDir(/etc/hive/conf) .setHadoopConfDir(/etc/hadoop/conf) .setDefaultDatabase(default) .build() );很多开发者在setHiveConfDir和setHadoopConfDir上踩坑这个路径必须是包含hive-site.xml和core-site.xml、hdfs-site.xml的目录不是单个文件路径。如果提交到YARN集群建议通过-Dyarn.ship-files把整个conf目录传上去否则本机能跑集群上就报Failed to create Hive Metastore client。defaultDatabase是默认库建表前不USE也能直接读写。我通常设成业务的ODS库省得每条SQL都写ods.xxx。4. 实操离线批读与实时写入的完整流程4.1 场景定义我拿一个实际改造过的网约车指标场景举例。业务表ride_order通过CDC进Kafka数据字段有订单号、乘客ID、司机ID、下单时间、上车点经纬度、订单状态、金额等。数仓里Hive有一张ODS层的ods_ride_order分区表分区字段是dt天还有一张DWS层的统计表需要按天统计每个司机完成订单数、总流水、平均抽成。传统做法是离线T1跑一遍HiveSQL实时报表从数据平台直接查Kafka的实时聚合。现在我们要做的是让Flink实时把Kafka数据写入ods_ride_orderHive离线继续消费这张表同时Flink还可以从Hive读取历史数据做回填。这套链路配合好实时和离线的口径就能统一。4.2 创建Hive Catalog并建表首先在Flink SQL客户端里创建CatalogCREATE CATALOG myhive WITH ( type hive, hive-conf-dir /etc/hive/conf, default-database default ); USE CATALOG myhive;然后在Hive侧直接建表也可以默认用Flink SQL建。建议表结构放到Hive侧管理这样离线任务也能用。表DDL如下CREATE TABLE ods_ride_order ( order_id BIGINT, passenger_id BIGINT, driver_id BIGINT, order_time TIMESTAMP, start_lng DOUBLE, start_lat DOUBLE, status INT, amount DECIMAL(10,2), ts TIMESTAMP ) PARTITIONED BY (dt STRING) STORED AS PARQUET TBLPROPERTIES ( sink.partition-commit.trigger partition-time, sink.partition-commit.delay 1 h, sink.partition-commit.watermark-time-zone Asia/Shanghai );注意我使用了TBLPROPERTIES这些是Flink Hive连接器识别的参数Hive本身不认识但会被Flink读取。设置partition-time触发提交并设置watermark时区是为了解决“跨天分区提交依赖水位线推进”的问题。4.3 批式读取Hive表Flink批模式读Hive表非常简单SET execution.runtime-mode BATCH; SELECT driver_id, COUNT(*) AS cnt, SUM(amount) AS total_amount FROM ods_ride_order WHERE dt 2024-12-01 GROUP BY driver_id;这里Flink底层会直接生成一个HiveTableSource把Hive表当成批式数据源读取。你可能会关心谓词下推比如WHERE dt 2024-12-01会不会只读对应分区目录答案是会的。Flink的Hive连接器会利用分区信息做分区裁剪也会把部分过滤条件下推到Hive读取阶段减少数据扫描量。但有个细节如果Hive表是ORC格式且用了Hive的ACID特性比如事务表Flink批读取会有问题可能会报不支持的事务类型。我建议数仓ODS层不用ACID表用普通外部表或者分区表保持文件组织简单Flink通用性最好。4.4 实时流写入Hive表含分区提交现在重点来了。实时从Kafka消费写入HiveFlink SQL的写法是CREATE TABLE kafka_source ( order_id BIGINT, passenger_id BIGINT, driver_id BIGINT, order_time TIMESTAMP(3), start_lng DOUBLE, start_lat DOUBLE, status INT, amount DECIMAL(10,2), ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH ( connector kafka, topic ride_order, properties.bootstrap.servers kafka:9092, properties.group.id flink-hive-sync, format json, json.ignore-parse-errors true ); INSERT INTO ods_ride_order SELECT order_id, passenger_id, driver_id, order_time, start_lng, start_lat, status, amount, ts, DATE_FORMAT(order_time, yyyy-MM-dd) AS dt FROM kafka_source;这条SQL会把每条数据按dt分区写入HDFS。写入文件是分块存储的Flink会按checkpoint间隔滚动文件。只有当checkpoint完成时文件的写入状态才算“提交”分区目录里才会出现完整文件。所以测试时如果写成UNCHECKED的批式也不对流式作业必须开启checkpointSET execution.checkpointing.interval 60s; SET execution.checkpointing.mode EXACTLY_ONCE;这里有个关键体验就算checkpoint开启分区提交的触发条件还没满足时你在Hive里查询该分区依然看不到数据。比如我们配置了partition-time加1 hour延迟一个12点产生的订单要等到13点多分区才正式可见。刚上手的时候多数人都会怀疑是不是写丢了实际上数据还在写入缓冲区里只是没提交。为了快速验证可以把触发方式临时改成process-time延迟设成0测试完再改回partition-time。提交触发参数如下sink.partition-commit.triggerprocess-time或partition-time。process-time以机器时间为准partition-time以分区时间和水位线为准。sink.partition-commit.delay延迟提交时间用于等待迟到数据。sink.partition-commit.policy.kind可以是metastore、success-file或两者都写。metastore表示把分区信息注册到MetaStoresuccess-file表示在分区目录写_SUCCESS文件。建议同时开启。生产环境我通常用sink.partition-commit.trigger partition-time, sink.partition-commit.delay 30min, sink.partition-commit.watermark-time-zone Asia/Shanghai, sink.partition-commit.policy.kind metastore,success-file这里的watermark-time-zone特别关键。Flink的partition-time提交是根据分区字段提取时间用当前水位线减去分区时间判断该分区是否到提交时间。如果集群默认时区是UTC你在中国时区跑分区16点的数据要被当成4点处理提交时间会足足晚8小时这种坑排查起来极其隐蔽。所以一定要显式设置水位线时区为Asia/Shanghai。4.5 流读Hive变更流读Hive是另一面的集成需求。比如离线任务每天往里写结果表Flink希望实时感知这些新数据用于下游实时计算。Flink Hive连接器支持以流模式读取Hive表的变化配置如下CREATE TABLE hive_stream_table ( driver_id BIGINT, cnt BIGINT, dt STRING ) WITH ( connector hive, streaming-source.enable true, streaming-source.partition-order partition-name, streaming-source.consume-start-offset 2024-12-01 ); SELECT * FROM hive_stream_table;这个功能的实现方式是Flink周期性扫描Hive分区目录发现新分区后把数据读入流中。所以表必须是分区表而且要保证新数据肯定写入新分区目录不能写到旧分区里否则流读无法感知。streaming-source.consume-start-offset用来指定从哪个分区开始消费不设就从头。监控间隔默认1分钟可以调streaming-source.monitor-interval。流读Hive听起来好用但生产上我不会作为主力实时链路。它的实时性受限于扫描周期一般是分钟级而且对HDFS文件可见性依赖较大跑批任务频繁改小文件时容易读到半成品文件。更稳妥的做法是让实时链路直接读Kafka用Hive流读做兜底或补充。5. 常见问题与排查实录5.1 依赖冲突NoSuchMethodError / ClassNotFoundException这类问题排在集成故障的第一位。症状五花八门作业启动时ClassNotFoundException: org.apache.hadoop.hive.conf.HiveConf运行中NoSuchMethodError: org.apache.hadoop.hive.metastore.api.Table.getParameters提交时报IncompatibleClassChangeError。我的排查套路先看Flink lib目录下有没有旧版本flink-connector-hive或其他hive相关jar有就清掉。再看自己项目的pom排除所有和Flink内置重复的依赖。尤其不要显式引入hadoop-client、hive-exec、hive-metastore除非你知道自己在做什么。可以使用mvn dependency:tree检查依赖关系锁定哪个包带进来了旧的guava或者protobuf。常见的坑是hive-exec依赖的guava版本和Flink冲突解决方法是排除掉hive-exec的guava依赖或者换用flink-sql-connector-hive来规避。经验之谈不要试图通过“多加一个jar”来修复依赖问题往往是越加越乱。把依赖树列出来找到冗余删掉才是正路。5.2 Hive分区表数据不更新我见过不少新人在测试“Flink实时写Hive”时用SHOW PARTITIONS看到分区存在但SELECT查不到数据。原因基本都是分区提交没完成或提交策略设置不对。排查顺序确认作业的checkpoint是否开启且正常完成。没有checkpointFlink写文件永远不会提交。查看HDFS上的分区目录有没有_SUCCESS文件。没有的话说明success-file策略没有触发或者还没到提交时间。检查partition-time对应的水位线是否推进。如果Kafka源没有定义水位线partition-time触发永远不生效。常见做法是给源表加WATERMARK FOR ts AS ts - INTERVAL 5 SECOND。检查时区设置。设置watermark-time-zone为本地时区并确认表上的分区字段格式是yyyy-MM-dd。一个小技巧在调试阶段用process-time触发延迟设0能立刻看到数据。确认链路通后再改回生产配置。5.3 HIVE小文件问题搜索热词里有“hive优化小文件”正好Flink写Hive更容易产生小文件。原因是Flink流式写入常按checkpoint间隔滚动文件checkpoint设得越短文件越多。如果一分钟一次checkpoint一天24小时会产生1440个分区文件这还只是一个分区的情况。数据量不大时可读性还行数据量一大NameNode压力直接爆炸Hive查询也会因为扫描小文件过多而慢到令人崩溃。应对思路有几个我会组合使用。思路一调大checkpoint间隔。把checkpoint从1分钟调到5分钟或10分钟文件数就能降到原来的五分之一到十分之一。但要注意checkpoint间隔越大故障恢复的延迟越高数据重复窗口变长。这个取舍要结合下游对实时性的容忍度。思路二Flink批式合并。定时启动一个Flink批作业读取前一天的分区用INSERT OVERWRITE重写一遍把分区目录里的小文件合并成大文件。Flink写文件时可以通过sink.partition-commit.policy.kind控制但合并操作本质是重写。我常用如下SQLSET execution.runtime-mode BATCH; INSERT OVERWRITE ods_ride_order SELECT ... FROM ods_ride_order WHERE dt 2024-12-01;这个操作就能把分区下的多个小文件重写成一个大的Parquet文件。前提是表是分区表、支持覆盖写并且没有ACID属性。思路三Hive侧合并。可以借助Hive自带的小文件合并参数比如hive.merge.mapredfilestrue、hive.merge.size.per.task256000000让跑批任务后自动合并。这依赖于Hive执行引擎对Flink写入的文件也有效但需要有人去触发合并作业。思路四减少分区粒度。如果不强依赖天级分区可以按小时或按周分区直接降低分区数量。但要评估查询特征不能为了减少分区而牺牲查询效率。我个人的生产做法是实时Flink按天分区写入凌晨3点启动一个Flink批合并作业把昨天的分区重写一次顺带清理异常数据。这样白天查询性能稳定实时性也不受影响。5.4 时区与事件时间问题除了前面说的提交时区Flink和Hive集成时时间字段的类型映射也容易出问题。Hive的TIMESTAMP映射到Flink是TIMESTAMP(9)但多数业务字段精度是毫秒或微秒读取时如果转换不当可能出现时间整体偏移。最简单的办法是在SQL里用CAST统一成毫秒级SELECT CAST(order_time AS TIMESTAMP(3)) FROM ods_ride_order;另外Hive的TIMESTAMP不带时区Flink的TIMESTAMP_LTZ带时区两者的语义完全不同。如果从Kafka读取的是order_time为字符串直接解析成TIMESTAMP_LTZ再写入Hive很可能出现小时级别的偏移。我在Kafka源定义时倾向于用TIMESTAMP(3)并且显式声明水位线然后写入Hive时再转成TIMESTAMP避免时区二次转换。5.5 性能调优实操笔记最后把这些年实践沉淀下来的几个调优点整理成表格方便对照。调优项推荐配置备注并发度sink.parallelism 分区数或较小值写入Hive的并发太高会生成大量小文件通常设为分区数即可批读取并行度分区并行度与文件数均衡避免启动过多task却只读一个文件Checkpoint间隔5~10分钟流写Hive文件大小、恢复时长、实时性三者平衡文件格式PARQUET Snappy压缩压缩比高Flink和Hive都原生支持文件大小控制128MB~256MB太大会影响Hive查询对HDFS友好内存taskmanager.memory.process.size4~8GB起步读Hive大表时bloom filter和orc相关缓存开销大读Hive谓词下推hive.java.opts无特殊要求主要确保HiveCatalog能正确获取分区信息实际操作中一个让人抓狂的隐形性能杀手是Hive表统计信息。Flink读取Hive时如果表的统计信息一直没更新优化器可能做出糟糕的执行计划比如把一个大表当成只有几条数据来优化。因此离线任务跑完或Flink写完分区后建议更新表的统计信息ANALYZE TABLE ods_ride_order PARTITION(dt2024-12-01) COMPUTE STATISTICS;如果表特别大至少保证分区级别的文件数和size是准确的。这个操作不能忘否则性能忽高忽低排查半天最后发现是统计信息过期。结尾关于这个方案的一点体会这些集成细节几乎都是我在“实时数仓和离线数仓口径统一”项目上一点一点啃出来的。Flink与Hive集成最大的价值不是让Flink能读写Hive这么简单而是让批和流的边界在存储层消融实时写入的数据离线任务能直接消费离线的历史数据Flink也能随时重新读取回填。真正做起来你会发现批流一体的难度不在计算引擎而在元数据、文件布局、提交语义和数据治理这些“脏活”上。最后再分享一个小技巧如果刚接触这套方案先别急着改造核心链路可以拿一张对实时性要求不高的维度表做试点用Flink实时写Hive跑一周观察文件数量、查询延迟和数据一致性。等这套机制跑稳了再逐步把实时明细、实时汇总迁到Hive上。批流一体不是一蹴而就的架构升级而是一个需要持续治理的数据工程实践。
返回列表