【发布时间】:2019-11-29 16:45:00
【问题描述】:
我想将 Spark 结构化流数据写入 cassandra。我的 spark 版本是 2.4.0。
我研究了一些帖子和一些使用 DataStax 企业平台。
我没用过,找到了一个方法foreachBatch,它有助于将流数据写入sink。
我查看了基于数据块 site 的文档。自己试试吧。
这是我写的代码:
parsed = parsed_opc \
.withWatermark("sourceTimeStamp", "10 minutes") \
.dropDuplicates(["id", "sourceTimeStamp"]) \
.groupBy(
window(parsed_opc.sourceTimeStamp, "4 seconds"),
parsed_opc.id
) \
.agg({"value": "avg"}) \
.withColumnRenamed("avg(value)", "avg")\
.withColumnRenamed("window", "sourceTime")
def writeToCassandra(writeDF, epochId):
writeDF.write \
.format("org.apache.spark.sql.cassandra")\
.mode('append')\
.options(table="opc", keyspace="poc")\
.save()
parsed.writeStream \
.foreachBatch(writeToCassandra) \
.outputMode("update") \
.start()
parsed 数据框的架构是:
root
|-- sourceTime: struct (nullable = false)
| |-- start: timestamp (nullable = true)
| |-- end: timestamp (nullable = true)
|-- id: string (nullable = true)
|-- avg: double (nullable = true)
我可以像这样成功地将这个流式 df 写入控制台:
query = parsed \
.writeStream \
.format("console")\
.outputMode("complete")\
.start()
控制台输出如下:
+--------------------+----+---+
| sourceTime| id|avg|
+--------------------+----+---+
|[2019-07-20 18:55...|Temp|2.0|
+--------------------+----+---+
所以,当写入控制台时,没关系。
但是当我在cqlsh 中查询时,表中没有附加任何记录。
这是 cassandra 中的表格创建脚本:
CREATE TABLE poc.opc ( id text, avg float,sourceTime timestamp PRIMARY KEY );
那么,你能告诉我有什么问题吗?
【问题讨论】:
-
您确定在尝试写入 Cassandra 时没有任何错误?
-
我不确定。但是程序运行时我没有在终端中看到任何错误。
-
我在 jupyter notebook 中运行这个程序。并且在终端中打印运行时日志。是否有任何查找日志文件的路径,可能是我丢失了?
标签: apache-spark cassandra pyspark spark-structured-streaming