【发布时间】: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 mode:https://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