【发布时间】:2019-06-08 22:43:46
【问题描述】:
我正在尝试从 kafka 主题中获取数据并将其推送到 hdfs 位置。我面临以下问题。
在每条消息 (kafka) 之后,hdfs 位置都会使用 .c000.csv 格式的部分文件进行更新。我在 HDFS 位置的顶部创建了一个配置单元表,但是 HIVE 无法读取从 spark 写入的任何数据结构化流媒体。
下面是spark结构化流处理后的文件格式
part-00001-abdda104-0ae2-4e8a-b2bd-3cb474081c87.c000.csv
这是我要插入的代码:
val kafkaDatademostr = spark.readStream.format("kafka").option("kafka.bootstrap.servers","ttt.tt.tt.tt.com:8092").option("subscribe","demostream").option("kafka.security.protocol","SASL_PLAINTEXT").load
val interval=kafkaDatademostr.select(col("value").cast("string")) .alias("csv").select("csv.*")
val interval2=interval.selectExpr("split(value,',')[0] as rog" ,"split(value,',')[1] as vol","split(value,',')[2] as agh","split(value,',')[3] as aght","split(value,',')[4] as asd")
// interval2.writeStream.outputMode("append").format("console").start()
interval2.writeStream.outputMode("append").partitionBy("rog").format("csv").trigger(Trigger.ProcessingTime("30 seconds")).option("path", "hdfs://vvv/apps/hive/warehouse/area.db/test_kafcsv/").start()
谁能帮帮我,为什么要创建这样的文件?
如果我这样做dfs -cat /part-00001-ad35a3b6-8485-47c8-b9d2-bab2f723d840.c000.csv,我可以看到我的值....但由于格式问题,它无法使用 hive 读取...
【问题讨论】:
-
好奇:你听说过 Kafka Connect 吗?你真的想要为这样一个简单的 Kafka 到 HDFS 用例编写 Spark 代码吗?另外,为什么 CSV 与 Parquet 相比,Hive 可以更好地阅读(不用担心引号和逗号)?
-
我的用例并不简单..在我的数据之上有一些 JOINS 和聚合...但我认为运行起来很简单..因为我是 spark/kafka 的初学者..所以我的简单用例也不起作用
-
好的,你能展示一下你的 Hive 表定义和你正在查看的数据吗?
-
我发现了一些东西....我们需要将最后一列指定为分区吗??
-
对于
INSERT INTO,是的。对于CREATE TABLE,您将改为使用PARTITIONED BY
标签: apache-spark hive apache-kafka spark-structured-streaming