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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 04:50:23