【问题标题】:is it possible to let spark structured stream(update mode) to write to db?是否可以让火花结构化流(更新模式)写入数据库?
【发布时间】:2021-01-03 01:47:10
【问题描述】:

我使用 spark(3.0.0) 结构化流从 kafka 读取主题。

我使用joins,然后使用mapGropusWithState 来获取我的流数据,所以我必须使用更新模式,根据我从spark官方指南:https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html#output-modes

spark 官方指南的下面部分没有提到DB sink,它也不支持files 写入update modehttps://spark.apache.org/docs/latest/structured-streaming-programming-guide.html#output-sinks

目前我将其输出到console,我想将数据存储在文件或数据库中。

所以我的问题是: 在我的情况下,如何将流数据写入数据库或文件? 我是否必须将数据写入 kafka,然后使用 kafka connect 将它们读回 files/db?

附言我按照文章获取aggregated流式查询。

- https://stackoverflow.com/questions/62738727/how-to-deduplicate-and-keep-latest-based-on-timestamp-field-in-spark-structured
- https://databricks.com/blog/2017/10/17/arbitrary-stateful-processing-in-apache-sparks-structured-streaming.html
- will also try one more time for below using java api
(https://stackoverflow.com/questions/50933606/spark-streaming-select-record-with-max-timestamp-for-each-id-in-dataframe-pysp)

【问题讨论】:

  • 不能使用 JDBC 写入数据库。 jdbcDF.write .format("jdbc") .option("url", "jdbc:postgresql:dbserver") .option("dbtable", "schema.tablename") .option("user", "username") 。选项(“密码”,“密码”).save()
  • 什么是数据库风格?
  • @Alex Ott,可以在内存一-H2

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


【解决方案1】:

我对 OUTPUT 和 WRITE 感到困惑。此外,我错误地假设 DB 和 FILE Sink 在文档的 OUTPUT SINK 部分中是并行的(因此在指南的 OUTPUT SINKs 部分中看不到 DB 接收器:https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html#output-sinks)。

我刚刚意识到 OUTPUT 模式(追加/更新/完成)是做查询流式查询约束。但这与如何写入 SINK 无关。我也意识到可以通过使用 FOREACH SINK 来实现 DB 写入(最初我只是理解它是为了额外的转换)。

我发现这些文章/讨论很有用

稍后,再次阅读官方指南,确认每个批次在写入存储时也可以执行自定义逻辑等。

【讨论】:

猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-05-31
  • 2019-11-26
  • 2020-02-25
  • 2016-09-14
  • 1970-01-01
  • 2020-06-21
相关资源
最近更新 更多