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

如何基于Smallrye Mutiny将返回List<String>的方法转为Multi<String>

改造方案

你当前实现的核心问题是先调用SAX的parse()方法把整个XML全量解析完成,所有customer节点全部堆在内存的rootNodes列表后,才开始遍历生成哈希推送数据,本质还是全量加载后批量返回,既没有实时性,大文件场景下还会占满内存。
SAX本身是事件驱动的解析模型,刚好和Mutiny Multi的推送逻辑匹配:不需要存所有解析完的节点,每解析完一个完整的<customer>标签,当场计算哈希值直接推送即可,解析完成后再通知流结束。


第一步:改造SaxHandler

去掉原来存储全量节点、全量哈希ID的列表,直接持有Mutiny的MultiEmitter,在单个customer节点解析完成的回调里当场算哈希推送,文档解析完成/出错时直接通知流的终态。顺便修复原代码里节点解析完没有回退父节点、digest变量未定义的bug:

import io.smallrye.mutiny.subscription.MultiEmitter;
import org.xml.sax.Attributes;
import org.xml.sax.helpers.DefaultHandler;
import javax.xml.bind.DatatypeConverter;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.util.HashMap;

public class SaxHandler extends DefaultHandler {
    private final MultiEmitter<? super String> emitter;
    private final String hashAlgorithm;
    private final HashMap<String, String> contextHeader = new HashMap<>();
    private ContextNode currentNode = null;
    private ContextNode rootNode = null;
    private final StringBuilder currentValue = new StringBuilder();

    public SaxHandler(MultiEmitter<? super String> emitter, String hashAlgorithm) {
        this.emitter = emitter;
        this.hashAlgorithm = hashAlgorithm;
    }

    @Override
    public void startElement(String uri, String localName, String qName, Attributes attributes) {
        if (rootNode == null && qName.equals("customer")) {
            rootNode = new ContextNode(contextHeader);
            currentNode = rootNode;
            rootNode.children.add(new ContextNode(rootNode, "type", qName));
        } else if (currentNode != null) {
            ContextNode n = new ContextNode(currentNode, qName, (String) null);
            currentNode.children.add(n);
            currentNode = n;
        }
    }

    @Override
    public void characters(char[] ch, int start, int length) {
        currentValue.append(ch, start, length);
    }

    @Override
    public void endElement(String uri, String localName, String qName) {
        if (rootNode != null && !qName.equals("customer")) {
            String value = currentValue.toString().trim().isEmpty() ? null : currentValue.toString().trim();
            currentNode.children.add(new ContextNode(currentNode, qName, value));
            // 回退到父节点,修复原代码层级错乱问题
            currentNode = currentNode.parent;
        }

        // 单个customer节点解析完成,当场计算哈希推送
        if (qName.equals("customer")) {
            try {
                String preHashString = rootNode.toString();
                MessageDigest digest = MessageDigest.getInstance(hashAlgorithm);
                byte[] hashBytes = digest.digest(preHashString.getBytes(StandardCharsets.UTF_8));
                String hashId = DatatypeConverter.printHexBinary(hashBytes).toLowerCase();
                emitter.emit(hashId);
            } catch (NoSuchAlgorithmException e) {
                emitter.fail(new RuntimeException("指定哈希算法不可用", e));
            }
            rootNode = null;
        }
        currentValue.setLength(0);
    }

    @Override
    public void endDocument() {
        // 整个XML解析完成,通知流结束
        emitter.complete();
    }

    @Override
    public void fatalError(org.xml.sax.SAXParseException e) {
        // 解析异常直接向下游传递
        emitter.fail(e);
    }
}

注意:你的ContextNode类需要保留parent成员变量,构造时正确赋值父节点引用即可。


第二步:改造主入口方法

不要先执行完解析再创建Multi,而是在Multi的发射器回调里启动SAX解析。因为SAX解析是阻塞操作,要把解析任务提交到专门的阻塞线程池执行,避免占满响应式事件循环线程,同时加取消逻辑:下游取消订阅时自动关闭输入流,停止多余解析。

import io.smallrye.mutiny.Multi;
import org.xml.sax.SAXException;
import javax.xml.parsers.ParserConfigurationException;
import javax.xml.parsers.SAXParserFactory;
import java.io.IOException;
import java.io.InputStream;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class HashGenerator {
    // 阻塞解析专用线程池,生产环境可根据并发量调整线程数
    private static final ExecutorService BLOCKING_IO_POOL = Executors.newCachedThreadPool();

    public Multi<String> xmlEventHashGenerator(InputStream xmlStream, String hashAlgorithm) {
        return Multi.createFrom().emitter(em -> {
            // 流终止时自动关闭输入流
            em.onTermination(() -> {
                try {
                    xmlStream.close();
                } catch (IOException ignored) {}
            });

            // 把阻塞解析任务丢到阻塞线程池执行
            BLOCKING_IO_POOL.submit(() -> {
                try {
                    SAXParserFactory factory = SAXParserFactory.newInstance();
                    factory.setFeature("http://apache.org/xml/features/disallow-doctype-decl", true);
                    SaxHandler saxHandler = new SaxHandler(em, hashAlgorithm);
                    factory.newSAXParser().parse(xmlStream, saxHandler);
                } catch (ParserConfigurationException | SAXException | IOException e) {
                    em.fail(e);
                }
            });
        });
    }
}

第三步:适配单元测试

Mutiny的Multi提供了阻塞等待的测试工具方法,不需要改原有测试逻辑,直接收集所有结果即可:

import org.junit.jupiter.api.Test;
import java.io.InputStream;
import java.util.List;

public class HashGeneratorTest {
    @Test
    public void xmlTest() {
        HashGenerator eventHashGenerator = new HashGenerator();
        InputStream xmlStream = getClass().getResourceAsStream("/customer.xml");
        List<String> eventHashIds = eventHashGenerator.xmlEventHashGenerator(xmlStream, "SHA3-256")
                .await().indefinitely()
                .collect().asList();
        System.out.println("\nGenerated Event Hash Ids : \n" + eventHashIds);
    }
}

改造后效果

  • 实时推送:每解析完一个customer节点就立刻计算哈希推给下游,不需要等整个文件解析完成
  • 内存友好:内存中只会保留当前正在解析的单个customer节点,哪怕GB级大XML文件也不会出现OOM
  • 异常透明:解析错误、算法不存在等异常会直接传递给下游,不会被吞掉只打印栈轨迹
  • 取消友好:下游如果中途取消订阅,会自动关闭输入流停止解析,不会做无用功
  • 兼容原有逻辑:如果需要拿全量结果,和原来一样调用collect方法收集即可,不需要改上层同步逻辑

内容的提问来源于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.29 03:48:22