Spark写入HDFS的CSV数据如何通过Hive查询及配置方法
var df = spark.readStream .format("kafka") //.option("kafka.bootstrap.servers", "kafka1:19092") .option("kafka.bootstrap.servers", "localhost:29092") .option("subscribe", "my_event") .option("includeHeaders", "true") .option("startingOffsets", "earliest") .load() df = df.selectExpr("CAST(topic AS STRING)", "CAST(partition AS STRING)", "CAST(offset AS STRING)", "CAST(value AS STRING)") val emp_schema = new StructType() .add("id", StringType, true) .add("timestamp", TimestampType, true) df = df.select( functions.col("topic"), functions.col("partition"), functions.col("offset"), functions.from_json(functions.col("value"), emp_schema).alias("data")) df = df.select("topic", "partition", "offset", "data.*") val query = df.writeStream .format("csv") .option("path", "hdfs://172.30.0.5:8020/test") .option("checkpointLocation", "checkpoint") .start() query.awaitTermination()
代码中hdfs://172.30.0.5:8020为HDFS的NameNode地址,当前Spark程序已可正常向该HDFS集群写入数据。
1. 如何通过Hive查询上述已写入HDFS的CSV格式数据?
通过建Hive外部表映射HDFS路径的方式即可查询,步骤如下:
- 登录Hive客户端,根据Spark写入的字段Schema建表,字段顺序、类型要和写入的CSV数据对齐:你写入的字段依次为
topic(STRING)、partition(STRING)、offset(STRING)、id(STRING)、timestamp(TIMESTAMP)。 - 推荐使用OpenCSVSerde处理CSV,避免字段内含逗号、引号导致解析错位,建表语句参考:
CREATE EXTERNAL TABLE IF NOT EXISTS my_event_kafka( topic STRING, partition STRING, offset STRING, id STRING, `timestamp` TIMESTAMP ) ROW FORMAT SERDE 'org.apache.hadoop.hive.serde2.OpenCSVSerde' WITH SERDEPROPERTIES ( "separatorChar" = ",", "quoteChar" = "\"", "escapeChar" = "\\" ) STORED AS TEXTFILE LOCATION 'hdfs://172.30.0.5:8020/test' TBLPROPERTIES ("skip.header.line.count"="0");
上述配置
skip.header.line.count=0是因为你的Spark写入代码没有开启header=true选项,CSV文件没有表头行;如果后续代码加了表头配置,把该值改为1即可。
- 建表完成后直接执行查询语句即可读取数据:
SELECT * FROM my_event_kafka LIMIT 10;
2. 是否必须将数据写入Hive可识别的特定目录才能被正常查询?
不需要。
Hive本身不强制要求数据存储在固定目录,只要Hive服务对目标存储路径有读权限,无论HDFS上的任意路径,还是其他兼容Hadoop API的存储路径,都可以通过建表时LOCATION参数指定路径映射,或者通过LOAD DATA命令将文件加载到表目录下完成查询。
需要注意的点:
- 路径下的文件格式要和建表时声明的存储格式一致
- 不要把Spark流式任务产生的checkpoint文件、以
_或.开头的临时文件放在数据目录下,Hive默认会忽略这类隐藏文件,若修改了默认配置扫描到这类文件,会触发解析报错
3. 是否必须为该数据存储目录创建对应的Hive数据库?具体操作流程是什么?
不需要。
Hive自带默认的default数据库,不单独新建数据库的话,直接在default库下建表映射路径即可正常查询。只有需要做业务数据隔离、权限管控的时候,才需要单独创建自定义数据库。
具体操作流程:
- (可选)新建自定义业务数据库,可指定数据库的默认存储根路径:
CREATE DATABASE IF NOT EXISTS stream_biz COMMENT '存储流式链路写入的业务数据' LOCATION 'hdfs://172.30.0.5:8020/user/hive/warehouse/stream_biz.db';
- 切换到目标数据库(也可以不切换,直接用
库名.表名的形式操作表):
USE stream_biz;
- 执行问题1中的建表语句,
LOCATION参数直接填写Spark写入的HDFS路径hdfs://172.30.0.5:8020/test即可,不需要把数据文件移动到数据库默认目录下。 - 验证:执行
SHOW TABLES;确认表存在,再执行抽样查询确认数据可正常解析。
注意:流式任务持续写入的场景下必须建外部表,不要建内部管理表,避免误删表时连带删除整个HDFS数据目录下的所有流式数据。
4. 代码中写入路径配置的test目录在文件系统中的实际存储位置是什么?
代码中配置的写入路径是完整HDFS URI:hdfs://172.30.0.5:8020/test,实际存储位置为HDFS集群根目录下的/test目录,对应NameNode服务地址为172.30.0.5:8020。
可以通过HDFS客户端命令查看目录下的文件:
# 写完整URI查看 hdfs dfs -ls hdfs://172.30.0.5:8020/test # 如果集群配置了默认FS为该NameNode地址,可简写为 hdfs dfs -ls /test
额外提示:你代码中checkpointLocation配置的是相对路径checkpoint,该路径会相对Spark任务提交时的工作目录解析,YARN集群模式下容易出现checkpoint丢失导致任务重复消费、启动失败的问题,建议改成完整的HDFS路径。
内容的提问来源于stack exchange,提问作者Sjoerd222888

