ARTICLE DETAIL

资讯详情

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

Flink与Hive函数生态整合:HiveModule复用与原生聚合加速实战

Flink与Hive函数生态整合:HiveModule复用与原生聚合加速实战 这几年做数据平台我见得太多了Hive 里沉淀了上百个自研 UDF/UDTF/UDAF从解析日志的正则函数到各种业务口径的聚合每一个都是踩过坑、对过账的资产。结果一上 Flink 实时计算很多团队第一反应就是“找人用 Java 重写一遍”一写就是两三个月中间还免不了在口径上跟离线侧扯皮。实际上 Flink 官方早就给了一条成熟的路——HiveModule 可以把 Hive 内置函数和元数据里的永久函数整体引入 Flink配合原生聚合加速让 sum/count/avg 这类高频聚合不必再经过 Hive 的 GenericUDAFEvaluator 包装层。这篇文章把我实际接入 Flink Hive 函数生态的经验拆开讲什么时候用 HiveModule、原生聚合加速怎么生效、UDF/UDTF/UDAF 各自的复用姿势、以及我在生产环境踩过的坑。1. 为什么非做不可存量函数资产和两条腿走路的现实1.1 存量资产Hive 函数库不是包袱而是金矿很多团队把 Hive UDF 当成历史包袱我反而觉得那是整个数仓里最值钱的资产之一。一个线上稳定跑了两年的 UDF意味着它已经经历过无数脏数据、边界值和口径变更的检验。你把它丢弃、重写表面上省了维护成本实际上是把这些隐性经验全部清零重新开始踩坑。我接过一个项目团队有 200 多个 Hive 函数涵盖埋点解析、IP 解析、用户画像标签加工、财务口径计算。当时实时链路要用其中大概 40 个如果全部用 Flink 原生ScalarFunction重写按一个人一天写 2 个函数算光开发就是 20 个工作日还不算单元测试、口径对齐、上线评审。最后我们走的是 HiveModule 复用三天就把 40 个函数全部接入实时作业跑通后面才逐步把高频函数替换成 Flink 原生实现。这并不是说重写没有价值而是说“先跑通、再优化”的顺序更符合生产现实。函数复用解决的是从 0 到 1 的问题原生重写解决的是从 1 到 10 的性能问题两者不冲突。1.2 一套 SQL 同时摸到实时表和离线表Flink 本身是流批一体的计算引擎加上 HiveCatalog 之后你可以在同一套 SQL 里既读 Kafka 的实时流又读 Hive 的离线分区表。函数层面也是一样同一个 Hive UDF既能用在实时流的字段解析上也能用在离线回扫和补数据的批任务里口径天然一致。我印象最深的一个场景是实时数仓的“首登用户”口径。离线侧这个口径写在一个 Hive UDAF 里用了两年业务部门已经认了这套逻辑。实时侧如果另写一套哪怕思路完全相同也总会有人质疑“两边是不是对不上”。直接用 HiveModule 加载同一个 UDAF至少在最开始的验证阶段能让两边跑出来的数字一模一样省掉大量扯皮时间。1.3 三个核心能力的分工与边界标题里的三件事其实对应 Flink Hive 集成的三条不同能力线HiveModule负责把 Hive 内置函数集注册进 Flink 的函数解析链同时让 Flink 能读到 Hive 元数据里的永久函数。原生聚合加速一个开关加一套原生实现让 count、sum、avg、min、max 这类常见内置聚合直接走 Flink 自己的 Aggregate 算子不经过 Hive 的 UDAF 评估器。函数复用通过HiveGenericUDF、HiveGenericUDTF、HiveGenericUDAF三个包装器把 Hive 侧的自定义函数整体搬到 Flink 里调用。边界也很清楚不是所有 Hive 内置函数都有原生实现自定义 UDF/UDTF/UDAF 永远走包装路径Hive UDAF 在无界流作业里要额外评估状态风险。搞清这三条边界后面遇到问题才不会慌。2. HiveModule 加载链路版本矩阵、依赖 Jar、解析优先级2.1 版本矩阵和两种加载姿势HiveModule 不是一个独立 Jar它包含在flink-connector-hive里。实际使用时有两条路SQL Client 用现成的 connector JarJava 项目用 Maven 依赖。先看 SQL Client 姿势。把对应版本的 connector 放进$FLINK_HOME/lib然后在 SQL 里加载模块# 以 Flink 1.18 Hive 3.1.3 为例 cp flink-sql-connector-hive-3.1.3_2.12-1.18.0.jar $FLINK_HOME/lib/LOAD MODULE hive WITH (hive-version 3.1.3); USE MODULES hive, core; CREATE CATALOG myhive WITH ( type hive, default-database default, hive-conf-dir /etc/hive/conf, hive-version 3.1.3 ); USE CATALOG myhive;Java 项目里对应这样写import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.TableEnvironment; import org.apache.flink.table.module.hive.HiveModule; import org.apache.flink.table.catalog.hive.HiveCatalog; EnvironmentSettings settings EnvironmentSettings.inStreamingMode(); TableEnvironment tEnv TableEnvironment.create(settings); // 1. 加载 HiveModule把 Hive 内置函数接入函数解析链 tEnv.loadModule(hive, new HiveModule(3.1.3)); // 2. 注册 HiveCatalog连接 Hive Metastore读取 Hive 表与永久函数 HiveCatalog catalog new HiveCatalog( myhive, default, /etc/hive/conf, 3.1.3); tEnv.registerCatalog(myhive, catalog); tEnv.useCatalog(myhive);Maven 依赖长这样dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-hive_2.12/artifactId version1.18.0/version /dependency dependency groupIdorg.apache.hive/groupId artifactIdhive-exec/artifactId version3.1.3/version /dependency版本兼容是这里最容易翻车的地方。Flink 对 Hive 的版本支持是有限列表不是随便配一个就能跑。常见稳定组合如下具体到你手里的 Flink 版本还是以官方文档的版本矩阵为准。Flink 版本官方支持的主要 Hive 版本我建议的主力组合1.15 ~ 1.171.2.3 / 2.1.1 / 2.3.6 / 3.1.32.3.6 或 3.1.31.18 ~ 1.202.3.6 / 3.1.33.1.3另一个容易忽略的点LOAD MODULE hive加载的是 Hive 内置函数它不需要连 Metastore 也能工作但你要用 HMS 里注册的永久函数就必须把 HiveCatalog 配好。很多人只加载了 Module 没配 Catalog然后发现自己的 UDF 找不到就是这个原因。2.2 函数解析优先级modules 就是个查字典的过程Flink 的函数解析机制可以理解成一本“字典”加载了多个模块就是多本字典叠在一起查函数时按顺序翻。默认情况下Flink 先查core模块Flink 内置函数再查其他模块。如果你执行了USE MODULES hive, core顺序就反过来了Hive 里的同名函数会优先命中。这个顺序不是小事。Hive 和 Flink 对某些同名函数的语义并不完全一致比如log这类带多参数顺序的函数、rand这种带随机状态的函数两边参数约定和实现细节都有差异。函数名一样不代表语义一样遇到怪结果先确认到底命中哪一边的实现。我的建议是默认保持USE MODULES core, hive让 Flink 内置函数优先只有当你明确知道某个函数必须用 Hive 的实现、且两边存在兼容性差异时才考虑调整顺序。顺序调错造成的不是报错而是静默的错误计算结果这类问题在生产上最致命。2.3 怎么确认加载成功加载完模块和 Catalog第一时间不要急着写业务 SQL先做两件事-- 1. 看 Hive 的函数有没有进来 SHOW FUNCTIONS LIKE %your_udf%; -- 2. 看函数归属和用法说明 DESCRIBE FUNCTION your_udf;如果函数没出现在结果里多半是 Module 没加载、Catalog 没连上 HMS或者函数根本不在这个库下面。在 Java 代码里可以用tEnv.listFunctions()做同样的检查。这一步 5 分钟就能排查完别等到作业提交了才发现函数不存在。3. 原生聚合加速到底加速了什么原理、开关、验证方法3.1 GenericUDAFEvaluator 的隐藏开销从哪来Hive 的自定义 UDAF 走的是GenericUDAFEvaluator那套评估模型一条数据从 Flink 进来要经过这么几步Flink 行数据转换成 Hive 的 inspect 对象调用iterate写入聚合 buffer最后terminate时再从 Hive 对象转回 Flink 类型。每一步都伴随 ObjectInspector 的类型推导、反射调用、临时对象分配。用大白话说Hive 的 UDAF 是为“离线批处理”设计的它假设数据可以慢慢处理每一行的转换开销无所谓。但放到 Flink 里尤其是大窗口聚合或者批式读取千万级分区数据时这部分转换开销会被放大得很明显。我测过一个自定义 UDAF同样的数据量走 Hive 包装路径比 Flink 原生AggregateFunction慢 3 到 4 倍而且 CPU 使用率明显偏高。3.2 table.hive.native-functions-enabled 的生效规则Flink 官方的解决方案是原生函数机制对应的配置项是table.hive.native-functions-enabled默认是开启的。这个开关的作用是在 SQL 优化阶段Flink 会优先把语义等价的 Hive 内置函数解析到自己的原生实现上。以聚合为例COUNT、SUM、AVG、MIN、MAX这类常用聚合会直接翻译成 Flink 的 Aggregate 算子不再实例化GenericUDAFEvaluator也不走 HiveObject 转换省掉的正是 3.1 节说的那几大开销。需要特别注意原生加速只覆盖“内置函数里语义等价的那部分”不是所有 Hive 内置函数都有原生版本。自定义 UDF/UDTF/UDAF 永远走包装路径不受这个开关影响。对于没有原生实现的函数Flink 会静默回退到 Hive 包装实现不会报错但性能就看你运气了。如果你发现某个内置聚合走了 Hive 包装路径想确认是不是开关被关了可以用SET table.hive.native-functions-enabled true;提示这个开关控制的是“函数解析”层面不是“是否允许连接 Hive”。关掉它不会让你连不上 Hive只是让所有函数都走 Hive 的实现逻辑。3.3 用 EXPLAIN 看计划用压测看收益判断一个聚合到底有没有命中原生加速别猜直接看执行计划EXPLAIN SELECT user_id, COUNT(*), SUM(amount) FROM myhive.default.orders GROUP BY user_id;如果计划里出现的是GroupAggregate、IncrementalGroupAggregate这类 Flink 原生算子就说明聚合走的是原生路径如果节点说明里带上了HiveGenericUDAF之类的字样那就是走了 Hive 包装路径。我在同一套资源下做过粗略压测数据全在内存、纯 CPU 计算4 个 TaskManager、每个 8 核场景对比如下场景原生聚合开启Hive 包装 UDAF2000 万行、10 个分组键、count/sum/avg约 9 秒约 28 秒5000 万行、字符串分组键、count约 21 秒约 63 秒这个数字只能说明趋势不同集群、不同数据分布会有差异但结论是稳定的能走原生就走原生省的是实打实的 CPU。生产上我建议把 EXPLAIN 检查纳入上线 checklist凡是涉及 Hive 聚合的作业至少确认一次执行计划里关键聚合没有被包成HiveGenericUDAF。4. 复用 Hive UDF/UDTF/UDAF 的三种姿势与差异化处理4.1 姿势一直接消费 HMS 里的永久函数这是我最推荐的方式。前提是你已经在 Hive 侧用CREATE FUNCTION把函数注册成了永久函数元数据存在 HMS 里。Flink 这边只要 HiveCatalog 配好、连得上 HMSSQL 里就能直接调用SELECT user_id, my_custom_udf(event_json) FROM myhive.default.event_log;这种方式的优点是零额外维护Hive 侧更新函数实现Flink 侧不用改任何 SQL 和 Jar。缺点是要注意函数 Jar 的可见性——HMS 里存的是函数名到类名的映射类文件本身还得放在能被 Flink 集群加载到的地方。最省事的做法是把函数 Jar 放到$FLINK_HOME/lib或者用ADD JAR明确声明ADD JAR hdfs:///udf-jars/my-udf.jar;4.2 姿势二Flink 侧临时注册如果你不想动 Hive 的元数据或者某些函数只在实时作业里临时用可以直接在 Flink 会话里注册临时函数ADD JAR hdfs:///udf-jars/my-udf.jar; CREATE TEMPORARY FUNCTION my_udf AS com.example.hive.MyUDF;临时函数生命周期只到会话结束适合验证、压测、临时口径。坏处是每次重启作业都要重新注册不适合长期生产任务。另外要注意用CREATE TEMPORARY FUNCTION注册的函数不会写进 HMSHive 那边是看不到的别搞混。4.3 UDF、UDTF、UDAF 的 SQL 姿势和差异三种函数类型在 Flink SQL 里的用法差别很大整理成一张表方便对照函数类型Flink SQL 用法底层包装实现要点UDFSELECT my_udf(col) FROM tHiveGenericUDF每行一次类型转换注意吞吐UDTFSELECT t.a FROM t, LATERAL TABLE(my_udtf(col1, col2)) AS t(a, b)HiveGenericUDTF必须给输出列起别名UDAFSELECT key, my_udaf(v) FROM t GROUP BY keyHiveGenericUDAF流式长期作业慎用UDTF 是最容易写错语法的地方。Hive 里是LATERAL VIEW explode(...)Flink 里要写成LATERAL TABLE(...) AS t(a, b)而且必须指定输出列的列名。我曾经见过同事漏写列别名直接报UDTFs result table needs alias查了半天才发现是语法问题。UDAF 在 Flink 里用起来最“省心”因为语法和普通聚合一模一样SELECT dim, my_custom_udaf(amount) AS total FROM detail_table GROUP BY dim;但省心不代表没风险UDAF 的状态问题在流式场景下是个大坑后面单独说。4.4 数据类型映射Hive 类型和 Flink 类型怎么对齐函数能正常调用底层依赖的是数据类型的互相转换。Hive 和 Flink 的类型系统不完全一样常见的映射关系如下Hive 类型Flink 类型备注TINYINT / SMALLINT / INT / BIGINT同名整数类型无损转换FLOAT / DOUBLEFLOAT / DOUBLE精度以 Hive 侧口径为准DECIMAL(p, s)DECIMAL(p, s)注意精度和标度对齐STRING / VARCHAR(n) / CHAR(n)STRINGCHAR 的尾部空格语义可能有差异BOOLEANBOOLEAN-BINARYBYTES-DATEDATE-TIMESTAMPTIMESTAMP(3) 为主高精度场景先做验证ARRAYTARRAYT嵌套类型有序列化开销MAPK, VMAPK, VKey 的类型不能太随意STRUCT...ROW...字段名大小写要特别注意UNIONTYPE不支持别在 Hive 表里用这个类型大部分转换是自动的但你心里得有数复杂类型ARRAY、MAP、ROW的转换开销比标量大得多。如果一个 UDF 的入参是STRUCT每行都要做一次反序列化吞吐一定上不去。碰到这种情况优先考虑在 SQL 层把复杂类型拆成标量列再传进去。5. 实战踩坑记录类冲突、函数解析和流批差异5.1 ClassNotFound / NoClassDefFound 的排查套路这是接 Hive 集成时碰到最多的报错。常见的有两种一种是ClassNotFoundException: org.apache.hadoop.hive.conf.HiveConf基本可以断定是hive-exec没进 classpath。SQL Client 场景检查 connector Jar 是否在lib目录Java 项目检查 Maven 依赖有没有把hive-exec打进去。另一种更隐蔽是类冲突。hive-exec里带了很多旧版本的第三方库比如旧版 Jackson、旧版 Thrift 和 Guava跟 Flink 自带的高版本 Guava 撞在一起表现就是各种诡异的NoSuchMethodError或序列化异常。这种情况优先用官方提供的flink-sql-connector-hive这个带 shade 的聚合 Jar它把冲突的依赖都重定位过了。Java 项目里如果坚持自己引hive-exec要做好 classloader 隔离的心理准备。我自己的经验是遇到hive相关的类加载问题先把taskmanager.classloader.resolve-order设置成child-first试一次。特别是用ADD JAR加载函数 Jar 的场景默认的 parent-first 策略可能让 TaskManager 优先从父加载器里找类你的函数类根本不会被看到。这条配置放在flink-conf.yaml里改了要重启集群才生效。5.2 同名函数覆盖和语义差异前面提到过模块顺序这里展开讲一个实际案例。我们有个作业要用 Hive 的rand做采样但默认模块顺序下 Flink 内置的rand先命中行为跟 Hive 版本不完全一致导致采样结果跟离线侧对不上。排查了大半天最后发现是函数解析顺序的问题。解决办法是在 SQL 里显式指定顺序USE MODULES hive, core;但这样做风险也很大因为 Hive 会把所有同名函数都抢占过来包括 Flink 内置的substring、concat这些。我的建议是尽量不要全局改顺序而是给容易混淆的函数起别名或者用 schema 级别的函数限定避免拿全局配置去赌局部需求。函数名一样不等于语义一样上线前先用一个小数据集对比两边的输出这一步不能省。5.3 Hive UDAF 在流式计算里的状态风险这是我认为最需要单独强调的坑。Hive 的GenericUDAFEvaluator把聚合中间态放在 evaluator 内部的 buffer 里Flink 包装层做的是把 Hive 的评估流程适配成 Flink 的聚合接口但 Hive 内部那部分状态能不能被 Flink 的 checkpoint 完整覆盖不是一个想当然的事情。我在无界流作业上实际遇到过 failover 之后聚合结果对不上的情况最后排查下来就是 Hive UDAF 的状态恢复存在问题。我的建议分三级批式作业有界流、离线回扫随便用Hive UDAF 没问题。流式短窗口作业可以先压测重点验证 checkpoint 恢复后的结果连续性。流式长期作业、全局聚合慎用。能用 Flink 原生AggregateFunction就用原生的性能更好状态也更可靠。这不是说 Hive UDAF 永远不能在流上用而是说它的状态模型和 Flink 的容错模型之间存在缝隙需要额外的验证成本。生产环境的稳定性优先级永远高于复用带来的研发效率。5.4 复杂类型和序列化性能的隐形损耗前文提过复杂类型有转换开销这里给一个可量化的例子。我测过一个入参为MAPSTRING, STRING的 Hive UDF在每秒 5 万条的数据流上单算子 CPU 使用率直接吃满一个核改成拆成两个STRING入参后CPU 降到 30% 左右。原因就是每行都要做一次 MAP 的反序列化和 ObjectInspector 转换。所以凡是能用标量表达的函数参数不要在 Hive 侧图省事传整个对象。这条建议同时适用于 UDF 和 UDTF。如果确实要传复杂类型至少做一层数据裁剪把不需要的字段在 SQL 层先过滤掉再传给函数。6. 生产落地的选型经验6.1 一张决策清单面对“这个函数要不要用 HiveModule 复用”这个问题我建议按下面的清单快筛场景建议理由离线批任务、有界流任务直接用 HiveModule HiveCatalog零重写、口径一致实时流式任务、标量/表函数先用 HiveModule 跑通再逐步替换先保证口径再优化性能实时流式任务、高频聚合优先 Flink 原生聚合性能和状态可靠性都更好函数实现简单、调用量大直接写 Flink 原生ScalarFunction一行函数重写收益稳定函数逻辑复杂、业务口径敏感保留 Hive 实现重写风险远大于性能收益长窗口、全局聚合、Kafka 输入避免 Hive UDAF状态恢复风险不可控6.2 一个混合落地的参考架构我最后落地这个项目时采用的是一个渐进式混合架构第一阶段离线所有函数继续留在 Hive实时作业全部通过 HiveModule 消费目标只是把链路跑通、口径对齐。第二阶段对实时作业里调用量前十的标量函数做性能分析把其中逻辑简单、未来大概率长期使用的替换成 Flink 原生ScalarFunction替换一个验证一个。第三阶段高频聚合全部走 Flink 原生聚合自定义 UDAF 仅保留在批式作业里。这套流程跑下来实时作业的吞吐比“全量 Hive 包装”方案提升了大约一倍但开发周期只比纯复用的方案多了一周。关键是每一步都有明确的验证点不会出现“全部推倒重写”的大爆炸。6.3 最后两手实操小技巧第一手上线前做一次函数解析巡检。把你作业里用到的所有函数列出来逐个跑EXPLAIN确认哪些走的是原生路径、哪些走了 Hive 包装路径。这个动作花不了半小时但能让你对作业的性能底牌心里有数。第二手建一个 SQL 回归用例集。每个要复用的函数准备一组固定输入分别在 Hive 里跑一遍、再在 Flink 里用 HiveModule 跑一遍结果必须完全一致。这个用例集平时不显眼但每次升级 Flink 版本、换 Hive 版本的时候它是你最后的防线。我吃过一次亏Flink 小版本升级后某个 Hive 内置函数的解析路径变了回归用例集第一时间发现了差异避免了一次生产事故。Flink 和 Hive 的函数生态整合核心思路就一句话离线资产要复用实时性能要原生。把这两个目标拆开按阶段推进就不会把好好的存量函数资产变成两套系统的历史包袱。
返回列表