You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Flink生成Map[String, Any]流并提取JSON指定字段

嘿,这三个Flink JSON处理的需求我都熟,给你一步步拆解解决:

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 06:09:12