【问题标题】:How to Write Structured Streaming Data into Cassandra with PySpark?如何使用 PySpark 将结构化流数据写入 Cassandra?
【发布时间】: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


【解决方案1】:

在研究主题后,我找到了解决方案。

仔细查看终端日志,我发现有一个错误日志: com.datastax.spark.connector.types.TypeConversionException: Cannot convert object [2019-07-20 18:55:00.0,2019-07-20 18:55:04.0] of type class org.apache.spark.sql.catalyst.expressions.GenericRowWithSchema to java.util.Date.

这是因为,在 spark 中执行 window 操作时,它会在时间戳列上的架构中添加一个结构,在本例中为 sourceTimesourceTime 的架构如下所示:

sourceTime: struct (nullable = false)
 |    |-- start: timestamp (nullable = true)
 |    |-- end: timestamp (nullable = true)

但我已经在 cassandra 中创建了一个列,该列已经是 sourceTime,但它只需要一个时间戳值。如果查看错误,它会尝试发送 cassandra 表中不存在的 startend timeStamp 参数。

因此,从 parsed 数据框中选择此列解决了问题: cassandra_df = parsed.select("sourcetime.start", "avg", "sourcetime.end", "id").

【讨论】:

    猜你喜欢
    • 2020-06-16
    • 2020-02-06
    • 2018-10-06
    • 1970-01-01
    • 1970-01-01
    • 2020-10-01
    • 2019-10-31
    • 2017-12-20
    • 2023-04-06
    相关资源
    最近更新 更多