【发布时间】:2022-08-12 17:41:16
【问题描述】:
我想根据特定条件将数据与火花流匹配,并且我想将此数据写入 Kafka。通过将 unmatched 保持在一个状态下,该状态将在 hdfs 中最多保留 2 天的数据。每个新传入的数据都会尝试匹配此状态下的未匹配数据。如何使用此状态事件? (我正在使用 pyspark)
标签: python apache-spark pyspark stream state
我想根据特定条件将数据与火花流匹配,并且我想将此数据写入 Kafka。通过将 unmatched 保持在一个状态下,该状态将在 hdfs 中最多保留 2 天的数据。每个新传入的数据都会尝试匹配此状态下的未匹配数据。如何使用此状态事件? (我正在使用 pyspark)
标签: python apache-spark pyspark stream state
Pyspark doesn't support stateful implementation by default。
只有 Scala/Java API 在KeyValueGroupedDataSet 上使用mapGroupsWithState 函数具有此选项
但是您可以将 2 天的数据存储在其他地方(文件系统或一些无 sql 数据库),然后对于每个新传入的数据,您可以转到 nosql 数据库并获取相应的数据并执行剩余的工作。
【讨论】: