ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Flink HiveCatalog 完全指南:以 Hive Metastore 为元数据中枢管理 Flink 表

Flink HiveCatalog 完全指南:以 Hive Metastore 为元数据中枢管理 Flink 表 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载HiveCatalog 是 Flink 内置的唯一持久化 Catalog 实现它复用 Hive MetastoreHMS作为元数据存储中枢让 Flink 用户创建的 Kafka 表、ES 表、Hive 表等元对象可以跨会话持久复用无需在每个 SQL Session 中重复建表。本文以官方文档为主线结合本仓库 flink-connector-hive 模块的源码实现系统讲解 HiveCatalog 的搭建、配置、两类表的语义差异、端到端实战示例以及 Flink/Hive 类型映射细节帮助你直接在现有 Hive 集群上开箱即用地管理 Flink 元数据。为什么需要 HiveCatalog在 Hadoop 生态中Hive Metastore 经过多年发展已成为事实上的元数据中心很多公司的生产环境只部署一个 HMS 服务实例用它统一管理 Hive 元数据乃至非 Hive 元数据作为唯一的 truth source。HiveCatalog正是基于这一现状提供两层价值Flink 与 Hive 共存对同时部署 Hive 和 Flink 的用户HiveCatalog允许直接使用现有的 HMS 来管理 Flink 的元数据实现两套引擎共享同一份元数据视图。纯 Flink 部署对只有 Flink 的用户HiveCatalog是 Flink 开箱即提供的唯一持久化 Catalog。没有持久化 Catalog 时用户只能通过 Flink SQL CREATE DDL 在每个会话中重复创建 Kafka 表等元对象浪费大量时间HiveCatalog填补了这一空白让表和其他元对象只需创建一次后续跨会话便捷引用和管理。从源码看HiveCatalog继承自 Flink 的AbstractCatalog基类实现了数据库、表、分区、函数、统计信息等全套 Catalog 接口底层通过HiveMetastoreClientWrapper封装对 HMS 的 Thrift 访问见 HiveCatalog.java。HiveCatalog设计为与现有 Hive 安装开箱即用兼容无需修改已有 HMS也无需改变表的数据布局或分区方式。搭建 HiveCatalog依赖准备设置HiveCatalog所需的依赖与整体 Flink-Hive 集成的依赖完全一致核心要点如下将 Hive 相关依赖加入 Flink 发行版的/lib/目录或在 Table API 程序中用-C、SQL Client 中用-l选项添加到 classpath通过HADOOP_CLASSPATH环境变量提供 Hadoop 依赖export HADOOP_CLASSPATHhadoop classpath推荐直接使用 Flink 预打包的flink-sql-connector-hive-2.3.10对应 Metastore 2.3.0 – 2.3.10或flink-sql-connector-hive-3.1.3对应 Metastore 3.0.0 – 3.1.3jar 包只有在预打包 jar 无法满足需求如 Hive 版本不在上述列表时才逐个添加hive-exec、libfb303、antlr-runtime等独立依赖。连接配置连接现有 Hive 安装需通过 Catalog 接口 和HiveCatalog完成支持 TableEnvironment 编程式注册与 YAML/SQL 声明式配置。以 SQL 方式为例CREATE CATALOG myhive WITH ( type hive, default-database mydatabase, hive-conf-dir /opt/hive-conf ); -- set the HiveCatalog as the current catalog of the session USE CATALOG myhive;在 Java 程序中则通过new HiveCatalog(name, defaultDatabase, hiveConfDir)构造后调用tableEnv.registerCatalog(myhive, hive)注册。创建HiveCatalog时支持以下选项对应源码 HiveCatalogFactoryOptions.java 中的ConfigOption定义Option必填默认值类型说明type是无StringCatalog 类型创建 HiveCatalog 时必须设为hivename是无StringCatalog 唯一名称仅 YAML 配置适用hive-conf-dir否无String包含hive-site.xml的 Hive 配置目录 URI。URI 需被 Hadoop FileSystem 支持相对 URI无 scheme视为本地文件系统路径。未指定时从 classpath 中查找hive-site.xmldefault-database否defaultStringCatalog 被设为当前 Catalog 时使用的默认数据库hive-version否无StringHiveCatalog能自动探测所用 Hive 版本不建议手动指定除非自动探测失败hadoop-conf-dir否无StringHadoop 配置目录路径仅支持本地文件系统路径。推荐用HADOOP_CONF_DIR环境变量设置仅当环境变量不可用如需为每个 Catalog 单独配置时才使用该选项从源码 HiveCatalogFactory.java 可以看到factoryIdentifier()返回hive创建时会把上述选项透传给HiveCatalog构造函数而 HiveCatalog.java 的createHiveConf方法会依次完成加载 Hadoop 配置core-site.xml/hdfs-site.xml/yarn-site.xml/mapred-site.xml→ 关闭 HiveConf 的静态配置加载 → 从hive-conf-dir指定的目录或 classpath读取hive-site.xml并合并进HiveConf。若指定的 Hadoop 目录下四种配置文件都不存在会直接抛出CatalogException。HiveCatalog 支持的两类表配置正确后HiveCatalog即可开箱即用用户通过 DDL 创建 Flink 元对象后能立刻看到。HiveCatalog可以处理两类表Hive 兼容表Hive-compatible tables这类表在元数据和存储层数据两方面都按 Hive 兼容的方式组织。因此通过 Flink 创建的 Hive 兼容表可以在 Hive 侧直接查询。官方推荐切换到 Hive dialect 来创建 Hive 兼容表如果使用默认 dialect则必须在表属性中设置connectorhive否则HiveCatalog默认将其视为通用表。使用 Hive dialect 时则无需设置connector属性。通用表Generic tables这类表是 Flink 特有的。用HiveCatalog创建通用表时只是借用 HMS 来持久化元数据。这些表虽然对 Hive 可见但 Hive 通常无法理解其元数据因此在 Hive 侧使用这类表会导致未定义行为。这一语义在源码中有清晰的体现HiveCatalog.java 中isHiveTable(Map)的实现是public static boolean isHiveTable(MapString, String properties) { return IDENTIFIER.equalsIgnoreCase(properties.get(CONNECTOR.key())); }即表属性中connector是否为hiveIDENTIFIER hive。对已存在的 Hive 表isHiveTable(Table)会读取 HMS 参数若存在已废弃的is_generic标记则按其取反判断否则只要表参数中不包含flink.connector/flink.connector.type前缀的属性即视为 Hive 表。相应地getCatalogTableType方法将表区分为HIVE_TABLE、FLINK_NON_MANAGED_TABLE与FLINK_MANAGED_TABLE三类且alterTable时不允许改变表的 Catalog 类型。端到端实战示例下面按官方文档的五个步骤完整走一遍从搭建 HMS 到用 Flink SQL 查询 Kafka 数据的过程。step 1搭建 Hive Metastore先让一个 Hive Metastore 运行起来。这里在本地启动 HMS并将hive-site.xml放在本地路径/opt/hive-conf/hive-site.xml配置内容如下configuration property namejavax.jdo.option.ConnectionURL/name valuejdbc:mysql://localhost/metastore?createDatabaseIfNotExisttrue/value descriptionmetadata is stored in a MySQL server/description /property property namejavax.jdo.option.ConnectionDriverName/name valuecom.mysql.jdbc.Driver/value descriptionMySQL JDBC driver class/description /property property namejavax.jdo.option.ConnectionUserName/name value.../value descriptionuser name for connecting to mysql server/description /property property namejavax.jdo.option.ConnectionPassword/name value.../value descriptionpassword for connecting to mysql server/description /property property namehive.metastore.uris/name valuethrift://localhost:9083/value descriptionIP address (or fully-qualified domain name) and port of the metastore host/description /property property namehive.metastore.schema.verification/name valuetrue/value /property /configuration用 Hive CLI 测试与 HMS 的连接执行一些命令可以看到当前有一个名为default的数据库且其中没有任何表hive show databases; OK default Time taken: 0.032 seconds, Fetched: 1 row(s) hive show tables; OK Time taken: 0.028 seconds, Fetched: 0 row(s)这里需要说明hive.metastore.uris指向的 Thrift 端口默认 9083正是HiveCatalog后续建立连接的入口。Flink 侧通过HiveMetastoreClientWrapper连接 HMS并在 HiveCatalog.open() 时校验default-database是否真实存在若不存在会直接抛出CatalogException这也是一个常见的排错点。step 2启动 SQL Client 并用 Flink SQL DDL 创建 Hive Catalog将全部 Hive 依赖加入 Flink 发行版的/lib目录然后在 Flink SQL CLI 中创建 Hive CatalogFlink SQL CREATE CATALOG myhive WITH ( type hive, hive-conf-dir /opt/hive-conf );hive-conf-dir指向 step 1 中存放hive-site.xml的目录。从源码 createHiveConf() 看该目录也可以是一个带 scheme 的 Hadoop FS URI如hdfs://...代码会通过Path.getFileSystem(hadoopConf)打开hive-site.xml并加载。step 3搭建 Kafka 集群启动一个本地 Kafka 集群创建一个名为test的 topic并向其中写入一些姓名 年龄的简单数据localhost$ bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test tom,15 john,21这些消息可以通过启动一个 Kafka 控制台消费者看到localhost$ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test --from-beginning tom,15 john,21step 4用 Flink SQL DDL 创建 Kafka 表用 Flink SQL DDL 创建一张简单的 Kafka 表并验证其 schemaFlink SQL USE CATALOG myhive; Flink SQL CREATE TABLE mykafka (name String, age Int) WITH ( connector kafka, topic test, properties.bootstrap.servers localhost:9092, properties.group.id testGroup, scan.startup.mode earliest-offset, format csv ); [INFO] Table has been created. Flink SQL DESCRIBE mykafka; root |-- name: STRING |-- age: INT注意这张 Kafka 表没有设置connectorhive因此它是一张通用表——元数据被持久化到 HMS 中但 Hive 侧无法理解其语义。此时在 Hive CLI 中验证可以看到表对 Hive 是可见的hive show tables; OK mykafka Time taken: 0.038 seconds, Fetched: 1 row(s)step 5运行 Flink SQL 查询 Kafka 表在 Flink 集群standalone 或 yarn-session 均可的 SQL Client 中执行一条简单的 select 查询Flink SQL select * from mykafka;继续向 Kafka topic 中生产更多消息localhost$ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test --from-beginning tom,15 john,21 kitty,30 amy,24 kaiky,18现在你可以在 SQL Client 中看到 Flink 产出的查询结果SQL Query Result (Table) Refresh: 1 s Page: Last of 1 name age tom 15 john 21 kitty 30 amy 24 kaiky 18至此mykafka表的元数据已经通过HiveCatalog持久化到 HMS。下次启动新的 SQL Client 会话时只需USE CATALOG myhive即可直接引用这张表无需重新执行CREATE TABLE。支持的数据类型与映射规则对于通用表HiveCatalog支持 Flink 的全部数据类型。对于 Hive 兼容表HiveCatalog需要将 Flink 数据类型映射为对应的 Hive 类型映射关系如下表Flink 数据类型Hive 数据类型CHAR(p)CHAR(p)VARCHAR(p)VARCHAR(p)STRINGSTRINGBOOLEANBOOLEANTINYINTTINYINTSMALLINTSMALLINTINTINTBIGINTLONGFLOATFLOATDOUBLEDOUBLEDECIMAL(p, s)DECIMAL(p, s)DATEDATETIMESTAMP(9)TIMESTAMPBYTESBINARYARRAYTLISTTMAPK, VMAPK, VROWSTRUCT类型映射的实现集中在 HiveTypeUtil.javaFlink → Hive 方向通过toHiveTypeInfo(DataType, checkPrecision)用LogicalTypeVisitor完成Hive → Flink 方向由toFlinkType(TypeInfo)按 Hive 类型类别PRIMITIVE/LIST/MAP/STRUCT递归转换其中 HiveTIMESTAMP统一映射为 FlinkTIMESTAMP(9)、HiveBINARY映射为BYTES。关于类型映射有几点需要注意Hive 的CHAR(p)最大长度为255Hive 的VARCHAR(p)最大长度为65535Hive 的MAP只支持基本类型的 key而 Flink 的MAP的 key 可以是任意数据类型Hive 的UNION类型不支持Hive 的TIMESTAMP精度恒为 9不支持其他精度不过 Hive UDF 可以处理精度 ≤ 9 的TIMESTAMP值Hive 不支持 Flink 的TIMESTAMP_WITH_TIME_ZONE、TIMESTAMP_WITH_LOCAL_TIME_ZONE和MULTISETFlink 的INTERVAL类型尚无法映射到 Hive 的INTERVAL类型。源码层面还有两处对上述规则的直接印证见 HiveTypeUtil.java当 FlinkCHAR(p)长度超过 Hive 上限 255、或VARCHAR(p)长度超过 65535 时若checkPrecisiontrue会抛出CatalogException错误信息中明确给出支持的长度区间[1, 255]/[1, 65535]若checkPrecisionfalse如调用 Hive UDF 处理数据时则自动提升为STRING类型。另外Flink 的STRING在内部等价于VARCHAR(Integer.MAX_VALUE)映射时会被识别并统一转换为 Hive 的STRING。进阶使用建议与注意事项用 Hive dialect 创建 Hive 兼容表官方明确推荐在 Flink 中使用 Hive dialect 执行 DDL 来创建 Hive 表、视图、分区和函数此时无需手动设置connector属性而使用默认 dialect 创建 Hive 兼容表时必须显式设置connectorhive。Hive 版本自动探测HiveCatalog通过HiveShimLoader.getHiveVersion()自动探测当前 classpath 中的 Hive 版本并按版本加载对应的HiveShim以屏蔽 Hive 各版本的 API 差异。官方建议不要手动指定hive-version仅在自动探测失败时才需要显式配置。表的类型不可随意变更从源码 HiveCatalog.java 看alterTable时会调用disallowChangeCatalogTableType校验Hive 兼容表与通用表之间不允许通过 ALTER 相互转换。元数据归属使用HiveCatalog创建的通用表虽然对 Hive 可见但只在 Flink 侧有确定的语义切勿在 Hive 侧直接使用这类表否则行为未定义。小结HiveCatalog是 Flink 与 Hive 生态衔接的桥梁对内它是 Flink 唯一的开箱即用持久化 Catalog让元对象跨会话复用对外它复用 HMS 这一业界元数据事实标准支持 Hive 兼容表与通用表两种语义并提供了完整的 Flink/Hive 类型映射。结合 overview.md 中支持的 Hive 版本列表2.3.0 – 2.3.10、3.1.0 – 3.1.3与 HiveCatalog.java 的实现细节你可以快速在现有 Hive 集群上搭建起统一的元数据管理方案。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Apache Flink HiveCatalog 完全指南以 Hive Metastore 为持久化元数据中心管理 Flink 元数据Apache Flink HiveCatalog 完全指南以 Hive Metastore 为持久化元数据中心管理 Flink 元数据 导读 HiveCat大数据流处理批处理数据工程Flink 与 Hive 集成指南HiveCatalog 元数据管理与 Hive 表读写实战Flink 与 Hive 集成指南HiveCatalog 元数据管理与 Hive 表读写实战 Apache Hive 早已是数据仓库生态系统的核心它既是面向大数据流处理批处理数据工程Flink 与 Apache Hive 集成指南HiveCatalog 持久化目录与 Hive 表读写实战Flink 与 Apache Hive 集成指南HiveCatalog 持久化目录与 Hive 表读写实战 Apache Hive 早已不仅是数据仓库生态中的大数据流处理批处理数据工程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表