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对此的处理逻辑,询问是否缺失配置项或需执行额外操作以加载数据。
解决方案
核心问题分析
- 文件格式不匹配:Hive表配置的是Parquet格式的SerDe和输入输出格式,但Kafka Connect写入HDFS的是Avro文件,Hive无法直接识别解析。
- 分区元数据未同步: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
相关产品推荐
相关产品推荐

