【问题标题】:Structured Spark Streaming multiple writes结构化 Spark Streaming 多次写入
【发布时间】:2019-07-08 20:34:48
【问题描述】:

我正在使用要写入 kafka 主题和 hbase 的数据流。 对于 Kafka,我使用如下格式:

dataset.selectExpr("id as key", "to_json(struct(*)) as value")
        .writeStream.format("kafka")
        .option("kafka.bootstrap.servers", Settings.KAFKA_URL)
        .option("topic", Settings.KAFKA_TOPIC2)
        .option("checkpointLocation", "/usr/local/Cellar/zookeepertmp")
        .outputMode(OutputMode.Complete())
        .start()

然后对于 Hbase,我会这样做:

  dataset.writeStream.outputMode(OutputMode.Complete())
    .foreach(new ForeachWriter[Row] {
      override def process(r: Row): Unit = {
        //my logic
      }

      override def close(errorOrNull: Throwable): Unit = {}

      override def open(partitionId: Long, version: Long): Boolean = {
        true
      }
    }).start().awaitTermination()

这会按预期写入 Hbase,但并不总是写入 kafka 主题。我不确定为什么会这样。

【问题讨论】:

  • 并不总是写到 kafka 主题。“并不总是”是什么意思?
  • 有时两者都可以正常工作,有时kafka部分根本不起作用
  • stackoverflow.com/questions/45331883/… 看起来类似的问题,但我看不到解决方案
  • 您在 Spark 日志中看到任何错误吗?司机/执行者?
  • 我没有看到你在第一个查询中使用awaitTermination

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


【解决方案1】:

在火花中使用foreachBatch

如果您想将流式查询的输出写入多个位置,那么您可以简单地多次写入输出 DataFrame/Dataset。但是,每次写入尝试都可能导致重新计算输出数据(包括可能重新读取输入数据)。为避免重新计算,您应该缓存输出 DataFrame/Dataset,将其写入多个位置,然后取消缓存。这是一个大纲。

streamingDF.writeStream.foreachBatch { (batchDF: DataFrame, batchId: Long) =>
    batchDF.persist()
    batchDF.write.format(…).save(…) // location 1
    batchDF.write.format(…).save(…) // location 2
    batchDF.unpersist()
}

【讨论】:

猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-07-12
  • 2018-10-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多