【发布时间】:2020-07-14 15:29:46
【问题描述】:
我正在尝试使用以下代码从 Kafka 主题中读取数据:
object Main {
def main(args: Array[String]) {
val sparkSession = createSparkSession()
val df = sparkSession.readStream.format("kafka").option("kafka.bootstrap.servers", "localhost:9092").option("subscribe", "test").option("startingOffsets", "earliest").load()
val df1 = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
df1.writeStream.format("parquet").option("format","append").option("checkpointLocation", "/home/krishna/Downloads/kafka_2.12-2.0.0/delete").option("path", "/home/krishna/Downloads/kafka_2.12-2.0.0/abc").option("truncate", "false").outputMode("append").start()
}
}
当我使用以下行时:
df1.writeStream
.format("console")
.option("truncate","false")
.start()
.awaitTermination()
然后输出将显示在控制台上。
但问题是当我在代码行下面替换上面的行时:
df1.writeStream
.format("csv")
.option("format","append")
.option("checkpointLocation", "/home/krishna/Downloads/kafka_2.12-2.0.0/delete")
.option("path", "/home/krishna/Downloads/kafka_2.12-2.0.0/abc")
.option("truncate", "false")
.outputMode("append")
.start()
那么输出不会以 CSV 格式保存。仅创建 abc 文件夹,并在其中创建元数据文件夹,但其中没有 CSV 文件。
我无法理解如果 o/p 成功显示在控制台上,那么为什么它没有以 csv、parquet 或文本的形式保存在文件中。
示例输出:
------------------
| key | value |
------------------
| null | abc |
| null | 123 |
|-----------------
依赖:
<dependencies>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_2.12</artifactId>
<version>2.4.5</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_2.12</artifactId>
<version>2.4.5</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming_2.12</artifactId>
<version>2.4.5</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql-kafka-0-10_2.12</artifactId>
<version>2.4.5</version>
<scope>provided</scope>
</dependency>
</dependencies>
【问题讨论】:
-
.option("truncate", "false") 用于控制台接收器,您已将其用于 csv 接收器不是吗?
-
我也试过上面的代码没有使用这个选项。但上面的代码无法将数据保存为 CSV,也没有显示任何错误。
-
您使用 df 用于控制台,使用 df1 用于 csv 可能是这种差异,您检查了吗?
-
@RamGhadiyaram 我更新了上面的问题,它只是拼写错误
-
当我使用读写而不是 readstream 和 writestream 时。然后这段代码就起作用了。
标签: scala apache-spark apache-kafka spark-streaming spark-structured-streaming