
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Apache Hive 早已是数据仓库生态系统的核心它既是面向大数据分析与 ETL 场景的 SQL 引擎也是一个用于数据发现、定义与演化的数据管理平台。本指南聚焦 Apache Flink当前仓库 flink 项目与 Hive 的两层集成其一利用 Hive Metastore 作为持久化 Catalog通过HiveCatalog跨会话共享 Flink 元数据其二让 Flink 直接读写 Hive 表统一批流处理能力。读完本文你将掌握 Hive 依赖的配置方式、HiveCatalog的多种连接与参数调优手段并理解流式读 Hive、时态表 Join、动态分区写入等核心能力的底层机制。Flink 与 Hive 集成的两个层面Flink 与 Hive 的集成包含两个层面利用 Hive Metastore 作为持久化 Catalog用户可通过HiveCatalog将不同会话中的 Flink 元数据存储到 Hive Metastore 中。例如用户可以使用HiveCatalog将 Kafka 表或 Elasticsearch 表存储在 Hive Metastore 中并在后续 SQL 查询中重新使用它们。利用 Flink 读写 Hive 表Flink 可作为 Hive 批处理引擎的一个性能替代选择也可以连续读写 Hive 表中的数据以支撑实时数据仓库应用。HiveCatalog的设计提供了与 Hive 良好的兼容性用户可以开箱即用地访问其已有的 Hive 数仓无需修改现有的 Hive Metastore也不需要更改表的数据位置或分区。从源码结构看Hive 连接器的核心实现位于 flink-connectors/flink-connector-hive 模块其中HiveCatalog类见 HiveCatalog.java封装了对 Hive Metastore 的全部访问逻辑而 HiveOptions.java 则集中定义了连接器相关的全部配置项。支持的 Hive 版本Flink 支持以下的 Hive 版本2.3 系列2.3.0、2.3.1、2.3.2、2.3.3、2.3.4、2.3.5、2.3.6、2.3.7、2.3.8、2.3.9、2.3.103.1 系列3.1.0、3.1.1、3.1.2、3.1.3请注意某些功能是否可用取决于你使用的 Hive 版本这些限制并非由 Flink 引起Hive 内置函数在使用 Hive-2.3.0 及更高版本时支持。列约束PRIMARY KEY 和 NOT NULL在使用 Hive-3.1.0 及更高版本时支持。更改表的统计信息在使用 Hive-2.3.0 及更高版本时支持。DATE列统计信息在使用 Hive-2.3.0 及更高版本时支持。依赖项配置要与 Hive 集成你需要在 Flink 的/lib/目录中添加一些额外的依赖包以便通过 Table API 或 SQL Client 与 Hive 进行交互或者你可以将这些依赖项放在专用文件夹中并分别使用 Table API 程序或 SQL Client 的-C或-l选项将它们添加到 classpath 中。由于 Apache Hive 构建于 Hadoop 之上首先需要准备 Hadoop 的依赖最简单的做法是导出 Hadoop classpathexport HADOOP_CLASSPATHhadoop classpath添加 Hive 依赖项有两种方式使用 Flink 提供的 Hive Jar 包或分别添加每个所需的 jar 包。如果你使用的 Hive 版本尚未在 Flink 预置列表中第二种方式会更适合。官方建议优先使用 Flink 提供的 Hive jar 包仅在不能满足需求时再考虑分开添加 jar 包。使用 Flink 提供的 Hive jar下表列出了所有可用的 Hive jar你可以选择一个并放到 Flink 发行版的/lib/目录中Metastore 版本Maven 依赖SQL Client JAR2.3.0 - 2.3.10flink-sql-connector-hive-2.3.10flink-sql-connector-hive-2.3.10含 scala 版本后缀仅稳定版本提供3.0.0 - 3.1.3flink-sql-connector-hive-3.1.3flink-sql-connector-hive-3.1.3含 scala 版本后缀仅稳定版本提供这两个连接器的模块定义可以在仓库中直接查看flink-sql-connector-hive-2.3.10/pom.xml 与 flink-sql-connector-hive-3.1.3/pom.xml。用户定义的依赖项不同 Hive 主版本所需的依赖项如下Hive 2.3.4/flink-version /lib // Flinks Hive connector包含 flink-hadoop-compatibility 和 flink-orc jar flink-connector-hivescala_version-version.jar // Hive 依赖 hive-exec-2.3.4.jar // 如果需要使用 Hive 方言添加 antlr-runtime antlr-runtime-3.5.2.jarHive 3.1.0/flink-version /lib // Flinks Hive connector flink-connector-hivescala_version-version.jar // Hive 依赖 hive-exec-3.1.0.jar libfb303-0.9.3.jar // 某些版本中 libfb303 未打包进 hive-exec需要单独添加 // 如果需要使用 Hive 方言添加 antlr-runtime antlr-runtime-3.5.2.jarMaven 依赖如果你在构建自己的应用程序则需要在 pom.xml 中添加以下依赖项。注意这些依赖应在运行时提供而不要将它们打进已生成的 jar 文件中因此使用provided作用域!-- Flink Dependency -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-hivescala_version/artifactId versionflink_version/version scopeprovided/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-table-api-java-bridgescala_version/artifactId versionflink_version/version scopeprovided/scope /dependency !-- Hive Dependency -- dependency groupIdorg.apache.hive/groupId artifactIdhive-exec/artifactId version${hive.version}/version scopeprovided/scope /dependency从 flink-connector-hive/pom.xml 可以看出连接器模块本身将hive-metastore、hive-exec、Hadoop 相关组件以及flink-table-*系列均声明为provided或optional作用域并排除了大量传递依赖如 guava、protobuf、log4j 等这正是为什么运行时必须显式提供 Hive/Hadoop classpath 的原因同时也避免了与 Flink 自身依赖的冲突。连接到 Hive通过 TableEnvironment 或者 YAML 配置使用 Catalog 接口 和HiveCatalog可以连接到现有的 Hive 集群。HiveCatalog的 Java 构造函数见 HiveCatalog.java依次接收 catalog 名称、默认 database、Hive 配置目录并支持可选的 Hadoop 配置目录与 Hive 版本参数其内部通过createHiveConf方法从hive-site.xml构建HiveConf并通过HiveShimLoader.loadHiveShim(hiveVersion)按版本加载对应的 Hive Shim 以屏蔽各 Hive 版本的 API 差异。以下是通过不同方式连接 Hive 的示例。JavaEnvironmentSettings settings EnvironmentSettings.inStreamingMode(); TableEnvironment tableEnv TableEnvironment.create(settings); String name myhive; String defaultDatabase mydatabase; String hiveConfDir /opt/hive-conf; HiveCatalog hive new HiveCatalog(name, defaultDatabase, hiveConfDir); tableEnv.registerCatalog(myhive, hive); // 将 HiveCatalog 设为当前 session 的 catalog tableEnv.useCatalog(myhive);Scalaval settings EnvironmentSettings.inStreamingMode() val tableEnv TableEnvironment.create(settings) val name myhive val defaultDatabase mydatabase val hiveConfDir /opt/hive-conf val hive new HiveCatalog(name, defaultDatabase, hiveConfDir) tableEnv.registerCatalog(myhive, hive) // 将 HiveCatalog 设为当前 session 的 catalog tableEnv.useCatalog(myhive)Pythonfrom pyflink.table import * from pyflink.table.catalog import HiveCatalog settings EnvironmentSettings.in_batch_mode() t_env TableEnvironment.create(settings) catalog_name myhive default_database mydatabase hive_conf_dir /opt/hive-conf hive_catalog HiveCatalog(catalog_name, default_database, hive_conf_dir) t_env.register_catalog(myhive, hive_catalog) # 将 HiveCatalog 设为当前 session 的 catalog t_env.use_catalog(myhive)YAMLSQL Client 配置execution: ... current-catalog: myhive # 将 HiveCatalog 设为当前 session 的 catalog current-database: mydatabase catalogs: - name: myhive type: hive hive-conf-dir: /opt/hive-confSQLDDL 方式CREATE CATALOG myhive WITH ( type hive, default-database mydatabase, hive-conf-dir /opt/hive-conf ); -- 将 HiveCatalog 设为当前 session 的 catalog USE CATALOG myhive;HiveCatalog 参数说明下表列出了通过 YAML 文件或 DDL 定义HiveCatalog时所支持的参数。这些参数在 HiveCatalogFactory.java 中被逐一注册并透传给HiveCatalog的构造函数参数必选默认值类型描述type是无StringCatalog 的类型。创建 HiveCatalog 时该参数必须设置为hive。name是无StringCatalog 的名字。仅在使用 YAML file 时需要指定。hive-conf-dir否无String指向包含 hive-site.xml 目录的 URI。该 URI 必须是 Hadoop 文件系统所支持的类型。如果指定一个相对 URI不包含 scheme则默认为本地文件系统。如果该参数没有指定会在 classpath 下查找 hive-site.xml。default-database否defaultString当一个 catalog 被设为当前 catalog 时所使用的默认当前 database。hive-version否无StringHiveCatalog 能够自动检测使用的 Hive 版本。建议不要手动设置 Hive 版本除非自动检测机制失败。hadoop-conf-dir否无StringHadoop 配置文件目录的路径。目前仅支持本地文件系统路径。推荐使用HADOOP_CONF_DIR环境变量来指定 Hadoop 配置仅在环境变量不满足需求时例如希望为每个 HiveCatalog 单独设置 Hadoop 配置再考虑使用该参数。HiveCatalog 的核心使用场景跨会话持久化 Flink 元数据Hive Metastore 在 Hadoop 生态中已事实上演变为元数据中心许多公司在生产环境中维护单一 Hive Metastore 服务实例将所有元数据无论 Hive 元数据还是非 Hive 元数据作为唯一事实来源。对于同时部署了 Hive 与 Flink 的用户HiveCatalog允许使用 Hive Metastore 管理 Flink 的元数据对于仅部署了 Flink 的用户HiveCatalog是 Flink 开箱即提供的唯一持久化 Catalog。没有持久化 Catalog 时用户使用 Flink SQLCREATE DDL不得不在每个会话中反复创建 Kafka 表等元对象浪费大量时间HiveCatalog通过让用户只创建一次表及其他元对象即可在之后的多个会话中便捷地引用和管理它们填补了这一空缺。Hive 兼容表与通用表HiveCatalog可以处理两种表Hive 兼容表Hive-compatible tables以 Hive 兼容的方式存储无论元数据还是存储层的数据都是 Hive 可识别的。因此通过 Flink 创建的 Hive 兼容表可以从 Hive 侧查询。通用表Generic tables仅 Flink 特有。使用HiveCatalog创建通用表时只是借助 HMS 持久化元数据。虽然 Hive 侧能看到这些表但 Hive 很可能无法理解其元数据在 Hive 中使用这类表会导致未定义行为。建议切换到 Hive 方言 创建 Hive 兼容表。如果使用默认方言创建 Hive 兼容表必须在表属性中设置connectorhive否则HiveCatalog默认将其视为通用表。注意使用 Hive 方言时不需要connector属性。一个完整的端到端示例下面通过一个简单示例演示HiveCatalog的完整使用链路。步骤 1搭建本地 Hive Metastore准备一个运行中的 Hive Metastore将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)步骤 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 );步骤 3搭建 Kafka 集群启动本地 Kafka 集群并创建名为test的 topic向其中生产若干 (name, age) 形式的简单数据localhost$ bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test tom,15 john,21通过 Kafka console consumer 可以验证消息localhost$ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test --from-beginning tom,15 john,21步骤 4用 Flink SQL DDL 创建 Kafka 表Flink 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验证该表在 Hive CLI 中也可见hive show tables; OK mykafka Time taken: 0.038 seconds, Fetched: 1 row(s)步骤 5运行 Flink SQL 查询 Kafka 表在 Flink 集群standalone 或 yarn-session的 SQL Client 中执行查询Flink SQL select * from mykafka;向 Kafka topic 继续生产更多消息后Flink SQL Client 中的查询结果会持续刷新输出SQL Query Result (Table) Refresh: 1 s Page: Last of 1 name age tom 15 john 21 kitty 30 amy 24 kaiky 18支持的数据类型映射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关于类型映射需要注意Hive 的CHAR(p)最大长度为 255Hive 的VARCHAR(p)最大长度为 65535Hive 的MAP仅支持基本类型作为 key而 Flink 的MAP可以是任意数据类型Hive 的UNION类型不受支持Hive 的TIMESTAMP精度恒为 9不支持其他精度而 Hive UDF 可以处理精度 ≤ 9 的TIMESTAMP值Hive 不支持 Flink 的TIMESTAMP_WITH_TIME_ZONE、TIMESTAMP_WITH_LOCAL_TIME_ZONE和MULTISETFlink 的INTERVAL类型目前还不能映射到 Hive 的INTERVAL类型。DDL 与 DMLDDL在 Flink 中执行 DDL 操作 Hive 的表、视图、分区、函数等元数据时建议使用 Hive 方言。DMLFlink 支持 DML 写入 Hive 表详细内容请参考 读写 Hive 表。读写 Hive 表的能力概览借助HiveCatalogFlink 可以对 Hive 表做统一的批和流处理。以下要点来自 读写 Hive 表 的完整说明便于你在配置时快速定位关键参数这些参数在 HiveOptions.java 中有源码级定义。流式读取 Hive 表Flink 支持以批和流两种模式从 Hive 表读取数据。批读基于执行查询时表的状态进行查询流读会持续监控表并在新数据可用时增量获取默认以批模式读取。流读支持分区表与非分区表分区表监控新分区的生成非分区表监控文件夹中的新文件。关键参数包括streaming-source.enable默认false是否启动流读。注意请确保每个分区/文件被原子地写入否则读取不到完整数据。streaming-source.partition.include默认all可选all与latest。latest读取按streaming-source.partition.order排序后的最新分区仅在流模式的 Hive 源表作为时态表时有效。streaming-source.monitor-interval默认None连续监控分区/文件的时间间隔。默认流读 Hive 间隔为1 min而流读 Hive 的 temporal join 默认间隔是60 min原因是当前实现中每个 TM 都要访问 Hive metastore可能对 metastore 产生压力。streaming-source.partition-order默认partition-name支持create-time比较文件系统修改时间、partition-time从分区名抽取时间和partition-name字典序比较。对于非分区表总是使用create-time。该选项与已弃用的streaming-source.consume-order等价。streaming-source.consume-start-offset默认None流模式起始消费偏移量。解析与比较方式取决于partition-order的设置。这些选项在 HiveOptions.java 中均有完整定义例如PartitionOrder枚举create-time/partition-time/partition-name明确说明了三种排序语义的差异。使用 SQL Hints 可以在不修改 Hive metastore 的情况下配置 Hive 表属性SELECT * FROM hive_table /* OPTIONS(streaming-source.enabletrue, streaming-source.consume-start-offset2020-05-20) */;注意事项监控策略是扫描当前位置路径中的所有目录/文件分区太多可能导致性能下降流读非分区表要求每个文件原子地写入目标目录流读分区表要求每个分区被原子地添加进 Hive metastore否则只有添加到现有分区的新数据会被消费流读 Hive 表不支持 Flink DDL 的 watermark 语法这些表不能被用于窗口算子。读取 Hive ViewsFlink 可以读取 Hive 中已定义的视图但存在两点限制Hive catalog 必须设置为当前 catalog 才能查询视图Table API 中使用tableEnv.useCatalog(...)SQL 客户端中使用USE CATALOG ...Hive 与 Flink SQL 语法不同关键字、字面值有差异请确保对视图的查询语法与 Flink 语法兼容。读取时的向量化优化与并行度推断当满足以下条件时Flink 会自动对 Hive 表进行向量化读取格式为 ORC 或 Parquet没有复杂类型列如 List、Map、Struct、Union。该特性默认开启可用以下配置禁用table.exec.hive.fallback-mapred-readertrue默认情况下Flink 会基于文件数量及每个文件中块的数量推断读取 Hive 的最佳并行度相关参数作用于当前作业所有 sourcetable.exec.hive.infer-source-parallelism.mode默认dynamic可选static作业创建阶段静态推断、dynamic作业执行阶段利用运行时信息更准确推断、none禁用。注意它仍受已弃用选项table.exec.hive.infer-source-parallelism影响需该值为true才能启用推断table.exec.hive.infer-source-parallelism.max默认1000source operator 推断的最大并发度默认值仅在静态推断模式下有效。时态表 Join你可以使用 Hive 表作为时态表让数据流通过 temporal join 关联 Hive 表。Flink 支持processing-time temporal join总是关联最新版本且支持分区表与非分区表对分区表可自动跟踪最新分区。注意Flink 还不支持 event-time temporal join Hive 表。最典型的场景是流作业中使用 Hive 表作为维度表Kafka 实时业务流/日志流通过 temporal join 关联每天由批任务更新的 Hive 维度表从而丰富数据流。示例见 读写 Hive 表 中的完整 SQL。对于最新分区作为时态表通过streaming-source.enable、streaming-source.partition.includelatest、streaming-source.monitor-interval、streaming-source.partition-order与partition.time-extractor.*组合配置对于最新表全量数据作为时态表则使用lookup.join.cache.ttl默认60 min控制缓存刷新周期。从 HiveOptions.java 可以看到lookup.join.cache.ttl的默认值为 60 分钟。写入 Hive 表Flink 支持批和流两种模式写入 Hive批模式作业完成时数据才可见支持追加INSERT INTO与覆盖INSERT OVERWRITE也支持静态分区与动态分区写入流模式持续添加新数据并通过分区提交使其可见不支持INSERT OVERWRITE。分区提交相关的典型配置来自 hive_read_write.mdSET table.sql-dialecthive; CREATE TABLE hive_table ( user_id STRING, order_amount DOUBLE ) PARTITIONED BY (dt STRING, hr STRING) STORED AS parquet TBLPROPERTIES ( partition.time-extractor.timestamp-pattern$dt $hr:00:00, sink.partition-commit.triggerpartition-time, sink.partition-commit.delay1 h, sink.partition-commit.policy.kindmetastore,success-file );注意如果在TIMESTAMP_LTZ列定义 watermark 并使用partition-time提交需要为sink.partition-commit.watermark-time-zone设置会话时区否则分区提交会延迟数小时。与写入相关的其他重要能力还包括table.exec.hive.fallback-mapred-writer控制是否使用 Flink 原生 writer设为false可对 parquet/orc 文件实现 S3 exactly-once 写入动态分区写入时默认按动态分区列排序table.exec.hive.sink.sort-by-dynamic-partition.enable默认true仅批模式生效自动收集统计信息table.exec.hive.sink.statistic-auto-gather.enable默认true仅批模式以及批/流模式下的小文件合并auto-compaction、compaction.small-files.avg-size默认 16MB、compaction.file-size、compaction.parallelism。已测试的文件格式Flink 的 Hive 集成已在以下文件格式上经过测试Text、CSV、SequenceFile、ORC、Parquet。结合 Hive 函数使用除了元数据与读写能力Hive 函数生态也可以直接复用详见 Hive FunctionsHiveModule将 Hive 内置函数作为 Flink 系统内置函数提供给 Flink SQL 与 Table API 用户可通过LOAD MODULE hive WITH (hive-version 2.3.4)加载原生 Hive 聚合函数从 Flink 1.17 起引入支持sum/count/avg/min/max五个函数通过table.exec.hive.native-agg-function.enabled默认false开启后可改用基于 hash 的聚合算子以显著提升聚合性能该选项在 HiveOptions.java 中有定义Hive UDF支持 UDF、GenericUDF、GenericUDTF、UDAF、GenericUDAFResolver2 等类型查询规划时自动翻译为 Flink 的 ScalarFunction、TableFunction 与 AggregateFunction。使用前提是包含该函数的 HiveCatalog 被设为当前 catalog且包含该函数的 jar 位于 Flink classpath。总结Flink 与 Hive 的集成是一套完整的元数据共享 统一批流计算方案HiveCatalog让 Hive Metastore 成为 Flink 元数据的持久化载体使 Kafka 表、Elasticsearch 表等元对象可以跨会话复用同时 Flink 原生读写 Hive 表的能力批读、流读、时态表 Join、分区提交、动态分区、文件合并、统计信息收集等使其成为实时数仓与离线数仓一体化建设的核心组件。配置时请牢记三点一是按 Hive 主版本选择正确的 connector jar 或手工依赖二是通过hive-conf-dir/hadoop-conf-dir/环境变量正确注入配置三是流读与分区提交相关参数需结合实际写入方式原子性谨慎设置。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐OpenMetadata与Hive集成大数据元数据管理终极指南在大数据时代企业面临着数据孤岛、元数据分散的严峻挑战。Hive作为企业级数据仓库的核心组件其元数据管理直接影响数据发现效率与协作能力。本文将为您展示如何通过数据目录数据血缘数据治理后端MCP 服务Omi Hive 集成插件实战用语音指令管理 Hive 项目、任务与搜索plugins/omi-hive-app 源码级解读Omi Hive 集成插件实战用语音指令管理 Hive 项目、任务与搜索plugins/omi hive app 源码级解读 本篇技术指南以开源仓库中 p人工智能AI 应用语音移动开发后端桌面应用智能硬件MCP 服务Flink CDC与Hive集成数据仓库的实时同步方案Flink CDC与Hive集成数据仓库的实时同步方案 1. 痛点与解决方案概述 传统数据仓库如Hive数据仓库依赖批量同步机制如Sqoop存在数据后端数据集成大数据流处理变更数据捕获数据同步上一篇php-code-coverage日志系统自定义日志处理器实现方法下一篇深入理解TikTokLive架构WebSocket通信与Protobuf协议解析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考