【发布时间】:2023-03-05 20:27:01
【问题描述】:
我有一个从 Kafka 读取的 spark 结构化蒸汽应用程序。 这是我的代码的基本结构。
我创建了 Spark 会话。
val spark = SparkSession
.builder
.appName("app_name")
.getOrCreate()
然后我从流中读取
val data_stream = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "server_list")
.option("subscribe", "topic")
.load()
在 Kafka 记录中,我将“值”转换为字符串。它从二进制转换为字符串。此时数据框中有1列
val df = data_stream
.select($"value".cast("string") as "json")
基于预定义的架构,我尝试将 JSON 结构解析为列。然而,这里的问题是如果数据是“坏的”,或者不同的格式,那么它与定义的模式不匹配。因此,下一个数据帧 (df2) 将空值放入列中。
val df2 = df.select(from_json($"json", schema) as "data")
.select("data.*")
我希望能够从 df2 中过滤掉在某个列(我用作数据库中的主键)中具有“null”的行,即忽略与架构不匹配的错误数据?
编辑:我在某种程度上能够做到这一点,但不是我想要的方式。
在我的流程中,我使用了一个使用.foreach(writer) 流程的查询。它的作用是打开与数据库的连接,处理每一行,然后关闭连接。 structured streaming 的文档提到了此过程所需的必需品。在 process 方法中,我从每一行获取值并检查我的主键是否为空,如果为空,我不将其插入数据库。
【问题讨论】:
标签: apache-spark apache-kafka spark-structured-streaming