Spark Java Structured Streaming多过滤条件处理及Kafka发送异常问题
Spark Java Structured Streaming多过滤条件处理及Kafka发送异常问题
看起来你现在遇到的问题很典型——因为你把两次过滤的结果都赋值给了同一个mainData变量,第二次的赋值直接覆盖了第一次的logout过滤结果,所以最后只有login的过滤数据被发送到Kafka啦。
给你整理了解决思路和具体的代码示例,帮你搞定这个问题:
首先先贴出你的原始参考数据,方便对照:
+-------+----------+------------+---------+---------------------+-----------+ |id |resource id|resource name|event-desc|event-date |ipaddress1 | +-------+----------+------------+---------+---------------------+-----------+ |2010001|119 |Netopia |logout |+56975-05-07 23:01:37|25:34:21:44| |2010001|119 |Netopia |login |+56975-05-07 23:01:37|25:34:21:44| |2010001|119 |Netopia |logout |+56975-05-07 23:01:37|25:34:21:45| |2010001|119 |Netopia |login |+56975-05-07 23:01:37|25:34:21:45| |2010001|119 |Netopia |logout |+56975-05-07 23:01:37|25:34:21:44| |2010001|119 |Netopia |logout |+56975-05-07 23:01:37|25:34:21:46| +-------+----------+------------+---------+---------------------+-----------+
核心解决思路
不要复用同一个变量存储不同过滤结果,为每个过滤条件创建独立的Dataset对象,这样两个结果都能被保留,再分别处理它们的Kafka发送逻辑。
修正后的Java代码示例
// 为不同过滤结果创建独立的Dataset,避免覆盖 Dataset<Row> logoutData = df.select("data.*").filter("data.`event-desc`='logout'"); Dataset<Row> loginData = df.select("data.*").filter("data.`event-desc`='login'"); // 处理logout数据发送到Kafka(Structured Streaming模式) StreamingQuery logoutQuery = logoutData .selectExpr("CAST(id AS STRING) AS key", "to_json(struct(*)) AS value") .writeStream() .format("kafka") .option("kafka.bootstrap.servers", "your-bootstrap-servers:9092") .option("topic", "your-logout-topic") // 可指定专属主题,也可与login共用同一主题 .option("checkpointLocation", "/path/to/unique-logout-checkpoint") .start(); // 处理login数据发送到Kafka StreamingQuery loginQuery = loginData .selectExpr("CAST(id AS STRING) AS key", "to_json(struct(*)) AS value") .writeStream() .format("kafka") .option("kafka.bootstrap.servers", "your-bootstrap-servers:9092") .option("topic", "your-login-topic") .option("checkpointLocation", "/path/to/unique-login-checkpoint") .start(); // 等待流查询运行,避免程序直接退出 logoutQuery.awaitTermination(); loginQuery.awaitTermination();
额外注意事项
- 字段名带横杠(如
event-desc)时,需要用反引号``包裹,避免Spark解析语法错误 - 每个流查询必须指定唯一的checkpoint路径,否则会出现元数据冲突
- 如果是批处理场景,把
writeStream替换为write,不需要checkpoint和awaitTermination - 若想将两种结果发送到同一个Kafka Topic,可先合并Dataset再发送:
Dataset<Row> combinedData = logoutData.union(loginData); // 后续执行合并数据的Kafka写入逻辑
备注:内容来源于stack exchange,提问作者Santosh Shimpi
相关产品推荐
相关产品推荐

