You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka Connect同步至HDFS Hive表无数据加载问题排查

问题描述

已配置Kafka Connect将Topic数据转储至启用Hive集成的HDFS,配置信息如下:

"confluent.topic.bootstrap.servers": "kafka-1:19092,kafka-2:29092,kafka-3:39092",
"connector.class": "io.confluent.connect.hdfs3.Hdfs3SinkConnector",
"flush.size": "3",
"hdfs.url": "hdfs://namenode:9000",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"logs.dir": "logs",
"name": "kafka to hdfs - repos",
"topics": "repos",
"topics.dir": "topics",
"value.converter": "io.confluent.connect.avro.AvroConverter",
"value.converter.schema.registry.url": "http://schema-registry:8081",
"hive.integration": "true",
"hive.metastore.uris": "thrift://hive-metastore:9083",
"schema.compatibility": "BACKWARD"

运行任务后,数据已写入HDFS:

# hdfs dfs -ls /topics/repos/partition=0
Found 3 items
-rw-r--r--   3 appuser supergroup        281 2022-12-02 13:42 /topics/repos/partition=0/repos+0+0000000000+0000000002.avro
-rw-r--r--   3 appuser supergroup        294 2022-12-02 13:42 /topics/repos/partition=0/repos+0+0000000003+0000000005.avro
-rw-r--r--   3 appuser supergroup        283 2022-12-02 13:42 /topics/repos/partition=0/repos+0+0000000006+0000000008.avro

Hive元数据中已创建对应表:

0: jdbc:hive2://localhost:10000> show tables;
+-----------+
| tab_name  |
+-----------+
| repos     |
+-----------+

但查询Hive表时无数据:

0: jdbc:hive2://localhost:10000> select * from repos;
+------------------+------------------+
| repos.repo_name  | repos.partition  |
+------------------+------------------+
+------------------+------------------+

查看表结构配置无明显异常:

0: jdbc:hive2://localhost:10000> DESCRIBE FORMATTED repos;
+-------------------------------+----------------------------------------------------+-----------------------------+
|           col_name            |                     data_type                      |           comment           |
+-------------------------------+----------------------------------------------------+-----------------------------+
| # col_name                    | data_type                                          | comment                     |
|                               | NULL                                               | NULL                        |
| repo_name                     | string                                             |                             |
|                               | NULL                                               | NULL                        |
| # Partition Information       | NULL                                               | NULL                        |
| # col_name                    | data_type                                          | comment                     |
|                               | NULL                                               | NULL                        |
| partition                     | string                                             |                             |
|                               | NULL                                               | NULL                        |
| # Detailed Table Information  | NULL                                               | NULL                        |
| Database:                     | default                                            | NULL                        |
| Owner:                        | null                                               | NULL                        |
| CreateTime:                   | Sat Dec 03 11:46:56 UTC 2022                       | NULL                        |
| LastAccessTime:               | UNKNOWN                                            | NULL                        |
| Retention:                    | 0                                                  | NULL                        |
| Location:                     | hdfs://namenode:9000/topics/repos                  | NULL                        |
| Table Type:                   | EXTERNAL_TABLE                                     | NULL                        |
| Table Parameters:             | NULL                                               | NULL                        |
|                               | COLUMN_STATS_ACCURATE                              | {"BASIC_STATS":"true"}  |
|                               | EXTERNAL                                           | TRUE                        |
|                               | bucketing_version                                  | 2                           |
|                               | numFiles                                           | 0                           |
|                               | numPartitions                                      | 0                           |
|                               | numRows                                            | 0                           |
|                               | rawDataSize                                        | 0                           |
|                               | totalSize                                          | 0                           |
|                               | transient_lastDdlTime                              | 1670068016                  |
|                               | NULL                                               | NULL                        |
| # Storage Information         | NULL                                               | NULL                        |
| SerDe Library:                | org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDe | NULL                        |
| InputFormat:                  | org.apache.hadoop.hive.ql.io.parquet.MapredParquetInputFormat | NULL                        |
| OutputFormat:                 | org.apache.hadoop.hive.ql.io.parquet.MapredParquetOutputFormat | NULL                        |
| Compressed:                   | No                                                 | NULL                        |
| Num Buckets:                  | -1                                                 | NULL                        |
| Bucket Columns:               | []                                                 | NULL                        |
| Sort Columns:                 | []                                                 | NULL                        |
| Storage Desc Params:          | NULL                                               | NULL                        |
|                               | serialization.format                               | 1                           |
+-------------------------------+----------------------------------------------------+-----------------------------+

日志中未发现错误,怀疑与partition=0子文件夹有关,不确定Hive对此的处理逻辑,询问是否缺失配置项或需执行额外操作以加载数据。

解决方案

核心问题分析

  1. 文件格式不匹配:Hive表配置的是Parquet格式的SerDe和输入输出格式,但Kafka Connect写入HDFS的是Avro文件,Hive无法直接识别解析。
  2. 分区元数据未同步:Hive表已定义partition分区字段,但Kafka Connect创建的partition=0目录对应的分区信息未注册到Hive元存储中,导致Hive无法发现该分区下的数据。

具体修复步骤

步骤1:修正Hive表的文件格式配置

将Hive表的SerDe、输入输出格式修改为适配Avro的配置,执行以下HiveQL语句:

ALTER TABLE repos SET SERDE 'org.apache.hadoop.hive.serde2.avro.AvroSerDe';
ALTER TABLE repos SET FILEFORMAT INPUTFORMAT 'org.apache.hadoop.hive.ql.io.avro.AvroContainerInputFormat' OUTPUTFORMAT 'org.apache.hadoop.hive.ql.io.avro.AvroContainerOutputFormat';

步骤2:同步Hive分区元数据

手动将partition=0分区添加到Hive表中:

ALTER TABLE repos ADD PARTITION (partition='0') LOCATION 'hdfs://namenode:9000/topics/repos/partition=0';

或者使用MSCK命令自动发现所有未注册的分区:

MSCK REPAIR TABLE repos;

步骤3:验证数据查询

再次执行查询语句,确认数据可以正常返回:

SELECT * FROM repos;

预防措施(可选)

如果希望Kafka Connect自动处理Hive分区和格式适配,可以补充以下Connector配置项:

  • 指定存储格式为Avro:"storage.class": "io.confluent.connect.hdfs3.storage.AvroStorage"
  • 关联分区字段:"hive.partition.field.name": "partition"

配置后,Connector会自动在写入数据时同步Hive分区元数据,并确保存储格式与Hive表兼容。

内容的提问来源于stack exchange,提问作者antontj

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.10 17:15:17