【问题标题】:spring-integration-aws dynamic file downloadspring-integration-aws 动态文件下载
【发布时间】:2018-06-09 13:14:17
【问题描述】:

我需要根据消息内容从 S3 下载文件。换句话说,要下载的文件以前是未知的,我必须在运行时搜索并找到它。 S3StreamingMessageSource 似乎不太合适,因为:

  1. 它依赖于轮询,因为我需要等待消息。
  2. 我找不到在流中间动态创建S3StreamingMessageSource 的任何方法。 gateway(IntegrationFlow) 看起来很有趣,但我需要的是一个不存在的 gateway(Function<Message<?>, IntegrationFlow>)。

另一个候选人是S3MessageHandler,但它不支持列出我需要查找所需文件的文件。

我可以直接使用 AWS API 实现我自己的消息处理程序,只是想知道我是否遗漏了什么,因为这似乎不是一个不寻常的要求。毕竟,并不是每个应用程序都只是坐在那里并不断轮询 S3 以获取新文件。

【问题讨论】:

  • 好吧,对于InputStream,您仍然可以使用S3RemoteFileTemplate 及其get() 函数或S3Session.readRaw(),如果您确实需要返回流。是的,我们可能会考虑将InputStream 支持添加到S3MessageHandler,但这不是错误
  • @ArtemBilan 我查看了get,但它调用了callback.doWithInputStream,然后关闭了流。与S3StreamingMessageSource 不同,消息没有机会向下游传输,任何需要做的事情都需要在回调中完成。我正在实现一个可以与S3StreamingMessageSource 一起使用的智能过滤器,这样我就不必重新发明轮子了。
  • 好的。我知道了。虽然这不是事件驱动的,但它仍然是可轮询的,但是您确实可以手动调用receive()。来自会话的readRaw() 应该对你有好处。

标签: spring amazon-s3 spring-integration enterprise-integration spring-integration-aws


【解决方案1】:

有S3RemoteFileTemplate 和list() 函数,您可以在handle() 中使用。然后split() 结果并调用S3MessageHandler 以下载每个远程文件。

虽然最后一个具有下载整个远程目录的功能。

【讨论】:

  • 我看了看,很遗憾它不起作用。 S3MessageHandler 无法流式传输文件,只能下载。这似乎是一个错误,因为如果有效负载不是File,它应该能够创建Message<InputStream>
  • 我成功了,请看我的回答。另见issue-82
【解决方案2】:

对于遇到这个问题的任何人,这就是我所做的。诀窍是:

  1. 稍后设置过滤器,而不是在构建时。注意没有addFilters或getFilters方法,所以过滤器只能设置一次,以后不能再添加。 @artem-bilan,这很不方便。
  2. 手动拨打S3StreamingMessageSource.receive。

    .handle(String.class, (fileName, h) -> {
    if (messageSource instanceof S3StreamingMessageSource) {
        S3StreamingMessageSource s3StreamingMessageSource = (S3StreamingMessageSource) messageSource;
    
        ChainFileListFilter<S3ObjectSummary> chainFileListFilter = new ChainFileListFilter<>();
        chainFileListFilter.addFilters(
                new S3SimplePatternFileListFilter("**/*/*.json.gz"),
                new S3PersistentAcceptOnceFileListFilter(metadataStore, ""),
                new S3FileListFilter(fileName)
        );
        s3StreamingMessageSource.setFilter(chainFileListFilter);
    
        return s3StreamingMessageSource.receive();
    }
    log.warn("Expected: {} but got: {}.",
            S3StreamingMessageSource.class.getName(), messageSource.getClass().getName());
    return messageSource.receive();
    }, spec -> spec
        .requiresReply(false) // in case all messages got filtered out
    )
    

【讨论】:

  • MessageSorce 并非设计为在运行时重新配置。当您同时访问它时,您可能会遇到一些竞争条件。但是您的addFilter 我们可以通过自定义ChainFileListFilter 实现来克服,您可以只将这个注入S3StreamingMessageSource 并在那里实现您自己的添加/删除逻辑。
  • 我认为您仍然可以自然地实现您的InputStream 要求,而无需调用MessageSource.receive() 手册。你需要一个S3RemoteFileTemplate 和它的getSession()。这个Session 有InputStream readRaw(String source) API。
  • @ArtemBilan 如果我使用Session,我将不得不自己进行过滤。我在S3RemoteFileTemplate 中没有看到任何过滤器。此外,这对我讨论的错误处理没有帮助 here
  • 好吧,您使用S3RemoteFileTemplate.list() 从.handle() 获取S3 资源列表,然后使用.filter() 拆分项目,然后提到Session.readRaw()。您可以将此流程包装到 .gateway() 中,并使用单点进行错误处理。
猜你喜欢
  • 2023-03-22
  • 1970-01-01
  • 2022-06-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多