如何使用SmallRye Mutiny Multi拼接不同类型的流?
解决SmallRye Mutiny拼接不同类型Multi的问题
核心原因
Multi.createBy().concatenating().streams要求所有待拼接的流元素类型必须一致,你的代码中一个是Multi<ByteArrayOutputStream>,另一个是Multi<Customer>,类型不匹配导致编译失败。
可行解决方案
方案1:统一到Object类型拼接
将两个流都转换为Multi<Object>,满足类型要求后再拼接,后续消费时通过类型判断区分处理:
修改你的generateCustomer方法:
public Multi<Object> generateCustomer(final Input input) { try { jsonGenerator.writeStartObject(); jsonGenerator.writeStringField("type", "customerDocument"); jsonGenerator.writeStringField("creationDate", Instant.now().toString()); jsonGenerator.writeFieldName("customerBody"); jsonGenerator.writeStartObject(); jsonGenerator.writeFieldName("customerList"); jsonGenerator.writeEndObject(); // End body jsonGenerator.writeEndObject(); // End whole json file // 确保JSON生成完成并刷新 jsonGenerator.flush(); Multi<Object> jsonWrapper = Multi.createFrom().item(jsonOutput); Multi<Object> customerStream = Multi.createFrom().publisher(new CustomerPublisher(input)) .map(customer -> (Object) customer); return Multi.createBy().concatenating().streams(jsonWrapper, customerStream); } catch (IOException ex) { throw new CustomerGeneratorException("Exception occurred during customer document wrapper creation : " + ex); } }
消费流时的类型判断逻辑:
generateCustomer(input) .subscribe().with(item -> { if (item instanceof ByteArrayOutputStream) { // 处理JSON输出流 ByteArrayOutputStream jsonStream = (ByteArrayOutputStream) item; String jsonContent = new String(jsonStream.toByteArray()); // 执行JSON相关业务逻辑 } else if (item instanceof Customer) { // 处理Customer对象 Customer customer = (Customer) item; // 执行Customer相关业务逻辑 } });
方案2:使用通用包装类封装元素
如果希望类型更安全,可以定义通用包装类统一元素类型,避免强制类型转换:
// 定义密封包装类,限定允许的子类类型 public sealed class Payload permits JsonPayload, CustomerPayload { } // 承载JSON输出流的包装类 public final class JsonPayload extends Payload { private final ByteArrayOutputStream jsonStream; public JsonPayload(ByteArrayOutputStream jsonStream) { this.jsonStream = jsonStream; } public ByteArrayOutputStream getJsonStream() { return jsonStream; } } // 承载Customer对象的包装类 public final class CustomerPayload extends Payload { private final Customer customer; public CustomerPayload(Customer customer) { this.customer = customer; } public Customer getCustomer() { return customer; } }
修改流转换与拼接逻辑:
public Multi<Payload> generateCustomer(final Input input) { try { jsonGenerator.writeStartObject(); jsonGenerator.writeStringField("type", "customerDocument"); jsonGenerator.writeStringField("creationDate", Instant.now().toString()); jsonGenerator.writeFieldName("customerBody"); jsonGenerator.writeStartObject(); jsonGenerator.writeFieldName("customerList"); jsonGenerator.writeEndObject(); // End body jsonGenerator.writeEndObject(); // End whole json file jsonGenerator.flush(); Multi<Payload> jsonWrapper = Multi.createFrom().item(new JsonPayload(jsonOutput)); Multi<Payload> customerStream = Multi.createFrom().publisher(new CustomerPublisher(input)) .map(CustomerPayload::new); return Multi.createBy().concatenating().streams(jsonWrapper, customerStream); } catch (IOException ex) { throw new CustomerGeneratorException("Exception occurred during customer document wrapper creation : " + ex); } }
消费时通过包装类类型区分处理:
generateCustomer(input) .subscribe().with(payload -> { if (payload instanceof JsonPayload) { ByteArrayOutputStream jsonStream = ((JsonPayload) payload).getJsonStream(); // 处理JSON内容 } else if (payload instanceof CustomerPayload) { Customer customer = ((CustomerPayload) payload).getCustomer(); // 处理Customer对象 } });
方案3:将JSON转换为Customer关联结构(业务允许时)
如果业务场景允许,可以把生成的JSON封装为一个特殊标记的Customer对象,直接用Multi<Customer>拼接:
public Multi<Customer> generateCustomer(final Input input) { try { jsonGenerator.writeStartObject(); jsonGenerator.writeStringField("type", "customerDocument"); jsonGenerator.writeStringField("creationDate", Instant.now().toString()); jsonGenerator.writeFieldName("customerBody"); jsonGenerator.writeStartObject(); jsonGenerator.writeFieldName("customerList"); jsonGenerator.writeEndObject(); // End body jsonGenerator.writeEndObject(); // End whole json file jsonGenerator.flush(); // 将JSON内容封装为特殊Customer对象(需根据你的Customer类结构调整字段) Customer headerCustomer = new Customer(); headerCustomer.setTag("DOCUMENT_HEADER"); headerCustomer.setExtContent(new String(jsonOutput.toByteArray())); Multi<Customer> jsonWrapper = Multi.createFrom().item(headerCustomer); Multi<Customer> customerStream = Multi.createFrom().publisher(new CustomerPublisher(input)); return Multi.createBy().concatenating().streams(jsonWrapper, customerStream); } catch (IOException ex) { throw new CustomerGeneratorException("Exception occurred during customer document wrapper creation : " + ex); } }
该方案要求Customer类有足够字段承载JSON内容,或业务允许存在此类特殊标记对象。
内容的提问来源于stack exchange,提问作者BATMAN_2008
相关产品推荐
相关产品推荐

