【问题标题】:Spark - Stream kafka to file that changes every day?Spark - 将 kafka 流式传输到每天都在变化的文件?
【发布时间】:2019-08-06 15:52:15
【问题描述】:

我有一个 kafka 流,我将在 spark 中处理。我想将此流的输出写入文件。但是,我想按天对这些文件进行分区,所以每天它都会开始写入一个新文件。可以做这样的事情吗?我希望它继续运行,当新的一天发生时,它将切换到写入新文件。

val streamInputDf = spark.readStream.format("kafka")
                    .option("kafka.bootstrapservers", "XXXX")
                    .option("subscribe", "XXXX")
                    .load()
val streamSelectDf = streamInputDf.select(...)

streamSelectDf.writeStream.format("parquet)
     .option("path", "xxx")
     ???

【问题讨论】:

  • 为什么不将Kafka中的数据直接消费到Spark中呢?
  • 因为出于审计目的,我们必须像我们拥有的所有其他数据流一样运行此数据流(每天在设定的时间)。所以我所要做的就是获取稍后将处理的数据。 @罗宾莫法特。我想过用.trigger(ProcessingTime("24 hours"))writeStream,但我不知道如何让文件被写入,真正改变
  • 你提到了 Kafka,但实际上 Kafka 本身应该只是充当消息总线。如果您使用的是 Cloudera/Hortonworks Data Flow 平台,您可以使用 NiFi 将数据移入/移出 Kafka,否则您可以使用 Spark 或 Kafka Connect 等工具来填补这个角色。
  • @DennisJaheruddin 正确。对不起,我忘了提。我将使用 spark 将这些数据处理到文件中。但是我不确定如何按日期对消息进行分区,以便在新的一天发生时将文件放入一个新文件中(对原始帖子的小编辑)

标签: apache-spark apache-kafka spark-streaming


【解决方案1】:

可以使用提供的partitionBy 从 spark 添加分区 DataFrameWriter 用于非流式传输或DataStreamWriter 用于 流数据。


以下是签名:

public DataFrameWriter partitionBy(scala.collection.Seq colNames)

DataStreamWriter partitionBy(scala.collection.Seq colNames) 按文件系统上的给定列对输出进行分区。

DataStreamWriter partitionBy(String... colNames) 分区 由文件系统上的给定列输出。

说明: partitionBy public DataStreamWriter partitionBy(String... colNames) 按文件系统上的给定列对输出进行分区。如果 指定时,输出布局在类似于 Hive 的文件系统上 分区方案。例如,当我们对数据集进行分区时 年和月,目录布局如下:

- 年=2016/月=01/ - 年=2016/月=02/

分区是最广泛使用的优化技术之一 物理数据布局。它提供了一个粗粒度的索引来跳过 当查询对分区有谓词时读取不必要的数据 列。为了使分区工作良好, 每列中的不同值通常应小于几十 数以千计。

参数:colNames - (undocumented) Returns: (undocumented) 由于: 2.0.0

所以如果你想按年和月对数据进行分区,spark 会将数据保存到如下文件夹:

年=2019/月=01/05 年=2019/月=02/05

选项 1(直接写入): 您提到了镶木地板 - 您可以使用以下方式保存为镶木地板格式:

df.write.partitionBy('year', 'month','day').format("parquet").save(path)

选项 2(使用相同的 partitionBy 插入配置单元):

你也可以像这样插入蜂巢表:

df.write.partitionBy('year', 'month', 'day').insertInto(String tableName)

获取所有 hive 分区:

Spark sql 是基于 Hive 查询语言的,所以你可以使用SHOW PARTITIONS

获取特定表中的分区列表。

sparkSession.sql("SHOW PARTITIONS partitionedHiveParquetTable")

结论: 我建议选项 2 ... 因为 Advantage 稍后您可以根据分区查询数据(也就是查询原始数据以了解您收到的内容),并且基础文件可以是 parquet 或 orc。

注意:

当您与SparkSessionBuilder 创建会话时,请确保您拥有.enableHiveSupport(),并确保您是否正确配置了hive-conf.xml 等。

【讨论】:

  • 如果您对答案没问题,请接受为所有者!
【解决方案2】:

基于this answerspark 应该能够根据年、月和日写入文件夹,这似乎正是您正在寻找的。还没有在火花流中尝试过,但希望这个例子能让你走上正轨:

df.write.partitionBy("year", "month", "day").format("parquet").save(outPath)

如果没有,您可以根据current_date() 放入一个变量文件路径

【讨论】:

    猜你喜欢
    • 2021-11-09
    • 2016-06-21
    • 2015-12-12
    • 1970-01-01
    • 2015-08-22
    • 1970-01-01
    • 2018-05-27
    • 2019-05-19
    • 2020-07-17
    相关资源
    最近更新 更多