【问题标题】:Not able to write csv file by taking data from kafka topic in spark Streaming with scala无法通过在使用 scala 的 spark Streaming 中从 kafka 主题中获取数据来编写 csv 文件
【发布时间】: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


【解决方案1】:

在控制台中,您使用的是df,而对于csv,您使用的是df1

大部分代码对我来说都很好。

试试这个。

df.writeStream 
    .format("csv")
    .option("format", "append")
    .trigger(processingTime = "5 seconds")
    .option("checkpointLocation", "/home/krishna/Downloads/kafka_2.12-2.0.0/delete")
.option("path", "/home/krishna/Downloads/kafka_2.12-2.0.0/abc")
    .outputMode("append")
    .start()

【讨论】:

    【解决方案2】:

    试试这个:

    df.writeStream
    .outputMode(OutputMode.Append())
    .format("csv")
    .option("checkpointLocation", "/home/krishna/Downloads/kafka_2.12-2.0.0/delete")
    .option("path", "/home/krishna/Downloads/kafka_2.12-2.0.0/abc/")
    .start()
    

    您可以使用格式类型:com.databricks.spark.csv

    【讨论】:

    • 与我发布的问题相同,不起作用
    【解决方案3】:

    我已经使用 Spark 2.4.5 测试了以下代码,它会根据需要生成 csv 文件:

     val sparkSession = SparkSession.builder()
        .appName("myAppName")
        .master("local[*]")
        .getOrCreate()
    
      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("csv")
        .outputMode("append")
        .option("set", ",")
        .option("checkpointLocation", "/home/krishna/Downloads/kafka_2.12-2.0.0/delete")
        .option("path", "/home/krishna/Downloads/kafka_2.12-2.0.0/abc.csv")
        .start()
        .awaitTermination()
    

    此代码将创建一个名为 abc.csv 的文件夹。根据 Sparksession 的并行性(由 spark.default.parallelism 配置),您会发现与分区一样多的 csv 文件。文件数反映了写出时 DataFrame 中的分区数。如果您在此之前对其进行重新分区,您最终会得到不同数量的文件。

    在我的情况下,分区是2,所以我在相应的文件夹中得到了这个输出:

    > ~/abc.csv$ ll
    total 28
    drwxrwxr-x  3 x x 4096 Apr 18 17:01 ./
    drwxr-xr-x 50 x x 4096 Apr 18 17:01 ../
    -rw-r--r--  1 x x    8 Apr 18 17:00 part-00000-77250d4a-e3af-46ef-b572-5476a3d075dd-c000.csv
    -rw-r--r--  1 x x    4 Apr 18 17:00 part-00000-82a76a8c-5977-4891-be36-1c2dc6837fb1-c000.csv
    

    【讨论】:

    • 当我使用读写而不是 readstream 和 writestream 时。然后这段代码就起作用了。
    • 代码没有问题,因为它在我的机器上运行良好。我猜你有依赖问题(可能还有冲突)。如前所述,我已经测试了代码并且运行良好。也许,您可以了解 what 不起作用?说“某事不工作”通常永远不会导致解决方案。
    • 正如我在问题标题中提到的,我的代码无法写入 CSV 文件。我也尝试了您的代码,它无法生成 CSV 文件。您可以在问题描述中看到依赖关系。
    • 再一次,如果您只是说“我的代码无法写入 CSV 文件”,但没有说明您遇到了什么错误/异常或您正在观察什么,则无法弄清楚您的问题是。
    • 另外,您如何提交/运行此 Spark 作业?看起来 pom 文件中有一些库被声明为“已提供”,因此您需要注意这一点。
    猜你喜欢
    • 2018-01-09
    • 2020-07-12
    • 1970-01-01
    • 1970-01-01
    • 2021-12-21
    • 2019-07-12
    • 2021-06-11
    • 2019-06-24
    • 2016-06-09
    相关资源
    最近更新 更多