【问题标题】:how to "group by" messages based on a header value如何根据标头值“分组”消息
【发布时间】:2019-08-14 14:11:49
【问题描述】:

我正在尝试根据遵循此标准的文件扩展名创建一个文件 zip:filename.{NUMBER},我正在做的是读取一个文件夹,按 .{number} 分组,然后创建一个唯一的以 .num 结尾的文件 .zip,例如:

文件夹/

文件.01

文件2.01

文件.02

文件2.02

文件夹 -> /已处理

file.01.zip 其中包含 -> file.01, file2.01

file02.zip 其中包含 -> file.02, file2.02

我所做的是使用 outboundGateway,拆分文件,丰富标题读取文件扩展名,然后聚合读取该标题,但似乎无法正常工作。

public IntegrationFlow integrationFlow() {
return flow
.handle(Ftp.outboundGateway(FTPServers.PC_LOCAL.getFactory(), AbstractRemoteFileOutboundGateway.Command.MGET, "payload")
                .fileExistsMode(FileExistsMode.REPLACE)
                .filterFunction(ftpFile -> {
                    int extensionIndex = ftpFile.getName().indexOf(".");
                    return extensionIndex != -1 && ftpFile.getName().substring(extensionIndex).matches("\\.([0-9]*)");
                })
                .localDirectory(new File("/tmp")))
            .split() //receiving an iterator, creates a message for each file
            .enrichHeaders(headerEnricherSpec -> headerEnricherSpec.headerExpression("warehouseId", "payload.getName().substring(payload.getName().indexOf('.') +1)"))
            .aggregate(aggregatorSpec -> aggregatorSpec.correlationExpression("headers['warehouseId']"))
            .transform(new ZipTransformer())
            .log(message -> {
                log.info(message.getHeaders().toString());
                return message;
            });
}

它给了我一条包含所有文件的消息,我应该期待 2 条消息。

【问题讨论】:

  • 我希望您什么也得不到,因为您还需要自定义发布策略。我建议第一步是打开 DEBUG 日志记录并遵循消息流。
  • 我去看看,我会告诉你的,谢谢!

标签: spring spring-integration spring-integration-dsl


【解决方案1】:

由于这个 dsl 的性质,我有一个动态数量的文件,所以我无法计算以相同数字结尾的消息(文件),我不认为超时可能是一个好的发布策略,我只是自己写了代码,没有写入磁盘:


.<List<File>, List<Message<ByteArrayOutputStream>>>transform(files -> {
                HashMap<String, ZipOutputStream> zipOutputStreamHashMap = new HashMap<>();
                HashMap<String, ByteArrayOutputStream> zipByteArrayMap = new HashMap<>();
                ArrayList<Message<ByteArrayOutputStream>> messageList = new ArrayList<>();
                files.forEach(file -> {
                    String warehouseId = file.getName().substring(file.getName().indexOf('.') + 1);
                    ZipOutputStream warehouseStream = zipOutputStreamHashMap.computeIfAbsent(warehouseId, s -> new ZipOutputStream(zipByteArrayMap.computeIfAbsent(s, s1 -> new ByteArrayOutputStream())));
                    try {
                        warehouseStream.putNextEntry(new ZipEntry(file.getName()));
                        FileInputStream inputStream = new FileInputStream(file);
                        byte[] bytes = new byte[4096];
                        int length;
                        while ((length = inputStream.read(bytes)) >= 0) {
                            warehouseStream.write(bytes, 0, length);
                        }
                        inputStream.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                });
                zipOutputStreamHashMap.forEach((s, zipOutputStream) -> {
                    try {
                        zipOutputStream.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                });
                zipByteArrayMap.forEach((key, byteArrayOutputStream) -> {
                    messageList.add(MessageBuilder.withPayload(byteArrayOutputStream).setHeader("warehouseId", key).build());
                });

                return messageList;
            })
            .split()
            .transform(ByteArrayOutputStream::toByteArray)
            .handle(Ftp.outboundAdapter(FTPServers.PC_LOCAL.getFactory(), FileExistsMode.REPLACE)
            ......

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-10-11
    • 1970-01-01
    • 2021-10-29
    • 2019-05-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多