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

如何为SmallRye Mutiny的Multi添加指定JSON包装器?

问题描述

我有一个Java方法,通过流方式生成JSON格式的客户信息,并用SmallRye Mutiny的Multi异步实时返回。现在想给生成的JSON添加指定结构的包装器,打算用Jackson JsonGenerator实现,觉得需要用Multi.createBy().concatenating()来组合逻辑,但不知道具体怎么把包装逻辑和现有返回的Multi结合起来。

现有方法代码

public static Multi<Customer> generateCustomer(final Input input) {
    try {
        return Multi.createFrom().publisher(new CustomerPublisher(input));
    } catch (Exception e) {
        throw new NewException("Exception occurred during the generation of Customer : " + e);
    }
}

当前返回的JSON结构

[
  {
    "name":"Batman",
    "age":45,
    "city":"gotham"
  },
  {
    "name":"superman",
    "age":50,
    "city":"moon"
  }
]

期望的JSON结构

{
  "isA": "customerDocument",
  "createdOn": "2022-10-10T12:29:43",
  "customerBody": {
    "customerList": [
      {
        "name": "Batman",
        "age": 45,
        "city": "gotham"
      },
      {
        "name": "superman",
        "age": 50,
        "city": "moon"
      }
    ]
  }
}

我尝试的代码(未成功结合)

public class CustomerGenerator {

    private ByteArrayOutputStream jsonOutput;
    private JsonGenerator jsonGenerator;

    private CustomerGenerator() {
        try {
            jsonOutput = new ByteArrayOutputStream();
            jsonGenerator = new JsonFactory().createGenerator(jsonOutput).useDefaultPrettyPrinter();
        } catch (IOException ex) {
            throw new TestDataGeneratorException("Exception occurred during the generation of customer document : " + ex);
        }
    }

    public static Multi < Customer > generateCustomer(final Input input) {
        CustomerGenerator customerGenerator = new CustomerGenerator();
        customerGenerator.wrapperStart();
        try {
            return Multi.createFrom().publisher(new CustomerPublisher(input));
        } catch (Exception e) {
            throw new NewException("Exception occurred during the generation of Customer : " + e);
        } finally {
            System.out.println("ALL DONE");
            customerGenerator.wrapperEnd();
        }
    }

    public void wrapperStart() {
        try {
            jsonGenerator.writeStartObject();
            jsonGenerator.writeStringField("type", "customerDocument");
            jsonGenerator.writeStringField("creationDate", Instant.now().toString());
            jsonGenerator.writeFieldName("customerBody");
            jsonGenerator.writeStartObject();
            jsonGenerator.writeFieldName("customerList");
        } catch (IOException ex) {
            throw new TestDataGeneratorException("Exception occurred during customer document wrapper creation : " + ex);
        }
    }

    public void wrapperEnd() {
        try {
            jsonGenerator.writeEndObject(); // End body
            jsonGenerator.writeEndObject(); // End whole json file
        } catch (IOException ex) {
            throw new TestDataGeneratorException("Exception occurred during customer document wrapper creation : " + ex);
        } finally {
            try {
                jsonGenerator.close();
                System.out.println("JSON DOCUMENT STRING : " + jsonOutput.toString());
            } catch (Exception e) {
                // do nothing
            }
        }
    }

}
解决方案

核心思路是将包装器开头、客户JSON流、包装器结尾三个部分作为Multi的独立流片段,通过Multi.createBy().concatenating()拼接,同时用Jackson流式生成JSON,避免内存堆积。

修改后的完整代码

import com.fasterxml.jackson.core.JsonFactory;
import com.fasterxml.jackson.core.JsonGenerator;
import com.fasterxml.jackson.databind.ObjectMapper;
import io.smallrye.mutiny.Multi;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.time.Instant;
import java.nio.ByteBuffer;

public class CustomerGenerator {

    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
    private static final JsonFactory JSON_FACTORY = OBJECT_MAPPER.getFactory();

    // 生成包装器开头的JSON字节片段
    private static Multi<ByteBuffer> generateWrapperStart() {
        ByteArrayOutputStream baos = new ByteArrayOutputStream();
        try (JsonGenerator generator = JSON_FACTORY.createGenerator(baos)) {
            generator.writeStartObject();
            generator.writeStringField("isA", "customerDocument");
            generator.writeStringField("createdOn", Instant.now().toString());
            generator.writeFieldName("customerBody");
            generator.writeStartObject();
            generator.writeFieldName("customerList");
            generator.writeStartArray(); // 开启customerList数组,之前的代码遗漏了这一步
        } catch (IOException e) {
            throw new TestDataGeneratorException("生成包装器开头失败", e);
        }
        return Multi.createFrom().item(ByteBuffer.wrap(baos.toByteArray()));
    }

    // 生成包装器结尾的JSON字节片段
    private static Multi<ByteBuffer> generateWrapperEnd() {
        ByteArrayOutputStream baos = new ByteArrayOutputStream();
        try (JsonGenerator generator = JSON_FACTORY.createGenerator(baos)) {
            generator.writeEndArray(); // 闭合customerList数组
            generator.writeEndObject(); // 闭合customerBody对象
            generator.writeEndObject(); // 闭合整个文档对象
        } catch (IOException e) {
            throw new TestDataGeneratorException("生成包装器结尾失败", e);
        }
        return Multi.createFrom().item(ByteBuffer.wrap(baos.toByteArray()));
    }

    // 将单个Customer对象序列化为JSON字节
    private static ByteBuffer serializeCustomer(Customer customer) {
        try {
            return ByteBuffer.wrap(OBJECT_MAPPER.writeValueAsBytes(customer));
        } catch (IOException e) {
            throw new TestDataGeneratorException("序列化Customer失败", e);
        }
    }

    public static Multi<ByteBuffer> generateCustomer(final Input input) {
        try {
            // 获取原始的Customer异步流
            Multi<Customer> customerMulti = Multi.createFrom().publisher(new CustomerPublisher(input));

            // 将Customer流转换为JSON字节流,并为非首个元素添加逗号分隔符
            Multi<ByteBuffer> customerJsonMulti = customerMulti
                    .map(CustomerGenerator::serializeCustomer)
                    .zipWithIndex()
                    .map(pair -> {
                        ByteBuffer customerBytes = pair.getItem1();
                        long index = pair.getItem2();

                        if (index > 0) {
                            // 非首个元素,前面添加逗号+空格
                            byte[] comma = ", ".getBytes();
                            ByteBuffer combined = ByteBuffer.allocate(comma.length + customerBytes.remaining());
                            combined.put(comma);
                            combined.put(customerBytes);
                            combined.flip();
                            return combined;
                        }
                        return customerBytes;
                    });

            // 拼接三个流:开头 -> 客户JSON流 -> 结尾
            return Multi.createBy().concatenating()
                    .streams(generateWrapperStart(), customerJsonMulti, generateWrapperEnd());
        } catch (Exception e) {
            throw new NewException("生成客户信息时发生异常", e);
        }
    }
}

关键说明

  1. 返回类型调整:将Multi<Customer>改为Multi<ByteBuffer>,适合HTTP流式响应场景,无需将所有数据加载到内存。
  2. 流式JSON生成:用JsonGenerator分别生成包装器的开头和结尾片段,确保JSON结构正确。
  3. 分隔符处理:通过zipWithIndex为非首个Customer的JSON添加逗号,保证数组格式合法。
  4. 流拼接:使用Multi.createBy().concatenating()将三个流按顺序拼接,实现异步实时输出完整包装后的JSON。

如果需要返回Multi<String>而非ByteBuffer,只需将ByteBuffer转换为字符串即可,例如:

.map(buffer -> new String(buffer.array()))

内容的提问来源于stack exchange,提问作者BATMAN_2008

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 22:01:38