如何为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); } } }
关键说明
- 返回类型调整:将
Multi<Customer>改为Multi<ByteBuffer>,适合HTTP流式响应场景,无需将所有数据加载到内存。 - 流式JSON生成:用
JsonGenerator分别生成包装器的开头和结尾片段,确保JSON结构正确。 - 分隔符处理:通过
zipWithIndex为非首个Customer的JSON添加逗号,保证数组格式合法。 - 流拼接:使用
Multi.createBy().concatenating()将三个流按顺序拼接,实现异步实时输出完整包装后的JSON。
如果需要返回Multi<String>而非ByteBuffer,只需将ByteBuffer转换为字符串即可,例如:
.map(buffer -> new String(buffer.array()))
内容的提问来源于stack exchange,提问作者BATMAN_2008
相关产品推荐
相关产品推荐

