如何基于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
相关产品推荐
相关产品推荐

