【问题标题】:Apache NiFi: Output to multiple FlowFiles simultaneously?Apache NiFi:同时输出到多个流文件?
【发布时间】:2016-11-13 02:40:22
【问题描述】:

有没有办法在 NiFi 的自定义处理器中同时写入不同的流?例如,我有第三方库使用类似这样的 API 进行重要处理:

public void process(InputStream in, OutputStream foo, OutputStream baa, List<String> args)
{
    ...
    foo.write(things);
    baa.write(stuff);
    ...
}

但我能找到的唯一示例都只使用一个输出流:

FlowFile transform = session.write(original, new OutputStreamCallback() {
        @Override
        public void process(OutputStream out) throws IOException {
            out.write("stuff");
        }
    });

处理是分批完成的(由于它的规模很大),因此执行所有处理然后写出单独的流程是不切实际的。

我能想到的唯一方法是多次处理输入:(

为了澄清,我想写入多个FlowFiles,使用session.write(flowfile, callback)方法,因此可以分别发送/管理不同的流

【问题讨论】:

  • 我不这么认为,我知道写入流文件的唯一方法是使用 OutputStreamCallback,它只有一个函数(进程),它只接受一个参数(一个 OutputStream)。
  • 是的,但是 TeeOutputStream 允许您拥有 1 个写入 2 个单独文件的流,不是吗(这还不够吗?
  • 我不相信,TeeOutputStream 将相同的东西写入两个流,我的函数没有(例如我的示例中的“事物”和“东西”)。谢谢
  • 另外,我想使用 session.write(flowfile, callback) 方法写入多个 flowfile (这不清楚,我会更新问题)。

标签: java apache-nifi


【解决方案1】:

NiFi API 基于一次作用于一个流文件,但您应该能够执行以下操作:

        FlowFile flowFile1 = session.create();
        final AtomicReference<FlowFile> holder = new AtomicReference<>(session.create());

        flowFile1 = session.write(flowFile1, new OutputStreamCallback() {
            @Override
            public void process(OutputStream out) throws IOException {

                FlowFile flowFile2 = session.write(holder.get(), new OutputStreamCallback() {
                    @Override
                    public void process(OutputStream out) throws IOException {

                    }
                });
                holder.set(flowFile2);

            }
        });

【讨论】:

  • 感谢这项工作。我唯一的改变是使第一个 OutputStream 具有不同的名称(和最终名称),这样您就可以在 inner-inner 函数中写入两者。
【解决方案2】:

由于您要从相同的输入产生不同的输出,您还可以考虑将这些步骤分解为专注于执行其特定功能的离散处理器。上面显示了“事物”和“东西”,例如,我建议您使用“DoThings”和“DoStuff”处理器。在您的流程中,您只需使用源连接两次即可向两者发送相同的流程文件。然后,这可以实现很好的并行操作,并允许它们具有不同的运行时/等。 NiFi 仍会为您维护出处线索,它实际上根本不会复制字节,而是将指针传递给原始内容。

【讨论】:

  • 我同意这有它的优点,但是如果每个输入流都完成了大量的处理,那么把它全部扔掉然后再做一次(或并行两次)感觉很浪费。它是一种权衡。不过还是谢谢。
  • 您能否在单个处理器中进行中间(或“共享”)处理,然后在操作不同时拆分为两个后续处理器?没有什么说当处理器 A 完成时流文件内容必须处于某种状态,只要 B 和 B' 都可以接受 A 的输出作为输入。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-10-14
  • 1970-01-01
  • 1970-01-01
  • 2020-03-27
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多