parquet 文件和 hive 表中的列名应该匹配,然后只有您可以使用您的 Hive 查询来查看特定列的数据。如果没有,您将看到这些列的值为 NULL 的行。
让我一步一步地告诉你它是如何写的:
1)创建一个包含列(id、name)的 Hive 表
0: jdbc:hive2://localhost:10000> CREATE EXTERNAL TABLE kmdb.test_ext_parquet (id STRING, name STRING) STORED AS PARQUET LOCATION '/km_hadoop/data/test_ext_parquet';
No rows affected (0.178 seconds)
0: jdbc:hive2://localhost:10000> SELECT * FROM kmdb.test_ext_parquet;
+----------------------+------------------------+--+
| test_ext_parquet.id | test_ext_parquet.name |
+----------------------+------------------------+--+
+----------------------+------------------------+--+
No rows selected (0.132 seconds)
2) 编写 Parquet 文件:我创建一个具有相同列名 (id,name) 的 spark 数据框并编写一个 parque 文件。
>>> dataDF=spark.createDataFrame([("1", "aaa"), ("2", "bbb")]) \
... .toDF("id", "name")
>>>
>>> dataDF.show(200, False)
+---+----+
|id |name|
+---+----+
|1 |aaa |
|2 |bbb |
+---+----+
>>> dataDF.coalesce(1).write.mode('append').parquet("/km_hadoop/data/test_ext_parquet")
>>>
3)验证数据
0: jdbc:hive2://localhost:10000> SELECT * FROM kmdb.test_ext_parquet;
+----------------------+------------------------+--+
| test_ext_parquet.id | test_ext_parquet.name |
+----------------------+------------------------+--+
| 1 | aaa |
| 2 | bbb |
+----------------------+------------------------+--+
2 rows selected (0.219 seconds)
0: jdbc:hive2://localhost:10000>
4)现在使用 Spark 数据框编写具有不同列名 (id2,name2) 的 parquet 文件
>>> dataDF2=spark.createDataFrame([("101", "ggg"), ("102", "hhh")]) \
... .toDF("id2", "name2")
>>>
>>> dataDF2.show(200, False)
+---+-----+
|id2|name2|
+---+-----+
|101|ggg |
|102|hhh |
+---+-----+
>>>
>>> dataDF2.coalesce(1).write.mode('append').parquet("/km_hadoop/data/test_ext_parquet")
>>>
5)让我们查询表来查看数据,查看列的 NULL 值。
0: jdbc:hive2://localhost:10000> SELECT * FROM kmdb.test_ext_parquet;
+----------------------+------------------------+--+
| test_ext_parquet.id | test_ext_parquet.name |
+----------------------+------------------------+--+
| 1 | aaa |
| NULL | NULL |
| 2 | bbb |
| NULL | NULL |
+----------------------+------------------------+--+
4 rows selected (0.256 seconds)
0: jdbc:hive2://localhost:10000>
6)现在,我将使用一个正确的列名 (id) 编写 parquet 文件,另一个不存在 (name2)。
>>> dataDF3=spark.createDataFrame([("201", "xxx"), ("202", "yyy")]) \
... .toDF("id", "name2")
>>>
>>> dataDF3.show(200, False)
+---+-----+
|id |name2|
+---+-----+
|201|xxx |
|202|yyy |
+---+-----+
>>>
>>> dataDF3.coalesce(1).write.mode('append').parquet("/km_hadoop/data/test_ext_parquet")
>>>
7)让我们查询表来查看数据,查看错误列的NULL值。
0: jdbc:hive2://localhost:10000> SELECT * FROM kmdb.test_ext_parquet;
+----------------------+------------------------+--+
| test_ext_parquet.id | test_ext_parquet.name |
+----------------------+------------------------+--+
| NULL | NULL |
| NULL | NULL |
| 1 | aaa |
| 2 | bbb |
| 201 | NULL |
| 202 | NULL |
+----------------------+------------------------+--+
6 rows selected (0.225 seconds)
0: jdbc:hive2://localhost:10000>
8) 现在我们如何读取“隐藏”数据?只需创建一个包含所有列的表。
0: jdbc:hive2://localhost:10000> CREATE EXTERNAL TABLE kmdb.test_ext_parquet_all_columns (id STRING, name STRING,id2 STRING, name2 STRING)
. . . . . . . . . . . . . . . .> STORED AS PARQUET LOCATION '/km_hadoop/data/test_ext_parquet';
No rows affected (0.097 seconds)
0: jdbc:hive2://localhost:10000> SELECT * FROM kmdb.test_ext_parquet_all_columns;
+----------------------------------+------------------------------------+-----------------------------------+-------------------------------------+--+
| test_ext_parquet_all_columns.id | test_ext_parquet_all_columns.name | test_ext_parquet_all_columns.id2 | test_ext_parquet_all_columns.name2 |
+----------------------------------+------------------------------------+-----------------------------------+-------------------------------------+--+
| NULL | NULL | 101 | ggg |
| NULL | NULL | 102 | hhh |
| 1 | aaa | NULL | NULL |
| 2 | bbb | NULL | NULL |
| 201 | NULL | NULL | xxx |
| 202 | NULL | NULL | yyy |
+----------------------------------+------------------------------------+-----------------------------------+-------------------------------------+--+
6 rows selected (0.176 seconds)
0: jdbc:hive2://localhost:10000>
现在,我如何知道 parquet 文件中的所有列?我可以通过使用 Spark 读取镶木地板文件来做到这一点:
>>> spark.read.option("mergeSchema", "true").parquet("/km_hadoop/data/test_ext_parquet").show(10,False)
+----+-----+----+----+
|id2 |name2|id |name|
+----+-----+----+----+
|101 |ggg |null|null|
|102 |hhh |null|null|
|null|xxx |201 |null|
|null|yyy |202 |null|
|null|null |1 |aaa |
|null|null |2 |bbb |
+----+-----+----+----+
>>>
我希望这会有所帮助..