【问题标题】:Create hourly snapshots with spark structured streaming Kafka使用 Spark 结构化流式 Kafka 创建每小时快照
【发布时间】:2020-07-19 07:23:59
【问题描述】:

我有一个用例来创建从 Kafka 主题消耗的数据的每小时快照。 我使用 spark 结构化流来使用来自 Kafka 的数据,并能够按照标准文档在控制台上打印它。 我是结构化流媒体的新手,不知道如何继续创建每小时快照。 有人可以帮我解决这个问题吗?

编辑:当前状态

假设来自kafka的事件的结构是

标识 ||姓名 ||位置

其中 id 唯一标识一条记录。 id 可以有多个事件,我想考虑最新事件。

假设我在上午 11:00 启动 spark 管道。 我想根据 Kafka 元数据汇总从晚上 11:00 到 12:00 到达的所有事件,并创建一个可以转储到文件的快照。

标识 ||姓名 ||位置

2  || ron  || 3.76
1  || can  || 2.68
4  || barn || 4.6

同样,我想汇总每小时的数据并生成当前小时的合并视图并将其转储到文件中。

我刚刚编写了将所有事件打印到控制台的基本消费者

val df = spark.readStream.format("kafka")
  .option("localhost:9092").option("subscribe", "geo-location").option("failOnDataLoss", true)
  .load.selectExpr("CAST(value AS STRING)")
  .as(Encoders.STRING)

val query = df.writeStream
  .format("console")
  .option("truncate", "false")

如果需要更多信息,请告诉我。 谢谢!

【问题讨论】:

  • 如何定义每小时快照?请显示一些输入和预期输出数据。另外,请分享您尝试过的方法以及似乎不起作用的部分。
  • 他们说不是免费的编码服务
  • 我在问题中添加了详细信息。
  • @raizsh 你觉得我的回答有帮助吗?
  • 我忙于另一项必须完成的任务。我会回到这个并更新线程。

标签: apache-kafka spark-structured-streaming


【解决方案1】:

您可以使用Trigger.ProcessingTime 并将其设置为1 小时。这将每隔一小时触发一次微批处理,从 Kafka 获取最新数据。

现在要使用此微批次创建快照,您需要对 id 进行重复数据删除,并根据数据中可用的时间戳字段选择最新记录。

df.writeStream
      .trigger(Trigger.ProcessingTime(1, TimeUnit.HOURS))
      .foreachBatch {
        (microBatch: DataFrame, batchId: Long) => {
          val snapshotDF = microBatch
            .withColumn("rnk",
              row_number().over(Window.partitionBy("id").orderBy(desc("timestampField")))
            ).filter("rnk = 1")

          // write snapshotDF to csv with timestamp appended to file name

        }
      }.start()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-01-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-11-28
    • 2019-07-29
    • 2019-10-03
    相关资源
    最近更新 更多