如何用Flink生成Map[String, Any]流并提取JSON指定字段
嘿,这三个Flink JSON处理的需求我都熟,给你一步步拆解解决:
处理Flink JSON DataStream的三个常见场景
1. 提取嵌套的User JSON结构
你的需求是把包含多层嵌套的原始JSON流,转换成只保留User对象的JSON流。这里用Flink内置的Jackson依赖(shaded版本,避免和项目其他Jackson依赖冲突)来处理最方便:
Java版本示例
import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.MapFunction; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper; public class FlinkJsonExtraction { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); ObjectMapper objectMapper = new ObjectMapper(); // 模拟原始输入JSON流 DataStream<String> rawJsonStream = env.fromElements( "{ \"Account Informations\": { \"User Info\": { \"User Name\": \"Albert Goldstein\", \"Date\": \"03/27/2015 08:35:11\", \"Location\": \"New York, USA\" }, \"User\": { \"Email\": \"FlinkIsDifficult@gmail.com\", \"Password\": \"*******\" } } }" ); // 提取并转换为目标JSON流 DataStream<String> userJsonStream = rawJsonStream.map(new MapFunction<String, String>() { @Override public String map(String rawJson) throws Exception { JsonNode rootNode = objectMapper.readTree(rawJson); // 逐层拿到嵌套的User节点 JsonNode userNode = rootNode.get("Account Informations").get("User"); // 序列化回JSON字符串 return objectMapper.writeValueAsString(userNode); } }); userJsonStream.print(); env.execute("Flink JSON Extraction Job"); } }
Scala版本示例
import org.apache.flink.streaming.api.scala._ import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper object FlinkJsonScalaExtraction { def main(args: Array[String]): Unit = { val env = StreamExecutionEnvironment.getExecutionEnvironment val objectMapper = new ObjectMapper() val rawJsonStream = env.fromElements( """{ "Account Informations": { "User Info": { "User Name": "Albert Goldstein", "Date": "03/27/2015 08:35:11", "Location": "New York, USA" }, "User": { "Email": "FlinkIsDifficult@gmail.com", "Password": "*******" } } }""" ) val userJsonStream = rawJsonStream.map { rawJson => val rootNode = objectMapper.readTree(rawJson) val userNode = rootNode.get("Account Informations").get("User") objectMapper.writeValueAsString(userNode) } userJsonStream.print() env.execute("Flink JSON Scala Job") } }
运行后就能得到你要的{"User": {"Email": "...", "Password": "..."}}格式的流。
2. 直接按字段解析获取指定值(比如Email)
必须支持!两种常用方式:
- 逐层解析后获取:拿到
User节点后直接调用get("Email").asText() - 用路径直接定位:用Jackson的
at()方法通过JSONPath直接获取,更简洁
示例代码(Java)
// 在map函数中添加这段逻辑 JsonNode rootNode = objectMapper.readTree(rawJson); // 方式1:逐层获取 String email1 = rootNode.get("Account Informations").get("User").get("Email").asText(); // 方式2:JSONPath直接定位 String email2 = rootNode.at("/Account Informations/User/Email").asText();
示例代码(Scala)
val rootNode = objectMapper.readTree(rawJson) val email = rootNode.at("/Account Informations/User/Email").asText()
两种方式都能直接拿到FlinkIsDifficult@gmail.com这个值。
3. 生成Map[String, Any]类型的流
不管是Java还是Scala都能轻松实现,这里给你两种思路:
方式一:从JsonNode自动转换为Map
利用Jackson的convertValue方法直接把User节点转成Map:
// Scala版本示例 val userMapStream = rawJsonStream.map { rawJson => val userNode = objectMapper.readTree(rawJson).get("Account Informations").get("User") import scala.collection.JavaConverters._ // 把Java Map转成Scala的Map[String, Any] objectMapper.convertValue(userNode, classOf[java.util.Map[String, Any]]).asScala.toMap }
方式二:手动构建Map(更灵活,适合只需要特定字段的场景)
// Scala版本示例 val userMapStream = rawJsonStream.map { rawJson => val userNode = objectMapper.readTree(rawJson).get("Account Informations").get("User") Map( "Email" -> userNode.get("Email").asText(), "Password" -> userNode.get("Password").asText() // 可以按需添加其他字段 ) }
Java版本的话,只需要把Scala的Map换成java.util.Map<String, Object>即可,逻辑完全一致。
内容的提问来源于stack exchange,提问作者TheEliteOne
相关产品推荐
相关产品推荐

