【问题标题】:Spark Structured Streaming with State (Pyspark)带状态的 Spark 结构化流式处理 (Pyspark)
【发布时间】:2022-08-12 17:41:16
【问题描述】:

我想根据特定条件将数据与火花流匹配,并且我想将此数据写入 Kafka。通过将 unmatched 保持在一个状态下,该状态将在 hdfs 中最多保留 2 天的数据。每个新传入的数据都会尝试匹配此状态下的未匹配数据。如何使用此状态事件? (我正在使用 pyspark)

    标签: python apache-spark pyspark stream state


    【解决方案1】:

    Pyspark doesn't support stateful implementation by default。

    只有 Scala/Java API 在KeyValueGroupedDataSet 上使用mapGroupsWithState 函数具有此选项

    但是您可以将 2 天的数据存储在其他地方(文件系统或一些无 sql 数据库),然后对于每个新传入的数据,您可以转到 nosql 数据库并获取相应的数据并执行剩余的工作。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-10-03
      • 2019-08-24
      • 1970-01-01
      • 1970-01-01
      • 2019-12-10
      • 1970-01-01
      相关资源
      最近更新 更多