【问题标题】:How can I write FlowFile attributes to Avro metadata inside the FlowFile's content?如何将 FlowFile 属性写入 FlowFile 内容中的 Avro 元数据?
【发布时间】:2018-01-29 20:42:10
【问题描述】:

我正在创建流文件,这些流文件在由ExecuteSql 处理器发出后在下游进行操作和拆分。我已经使用要放入每个 FlowFile 内容中包含的 Avro 元数据的数据填充了 FlowFiles 的属性。

我该怎么做?

我尝试使用配置了AvroReaderAvroRecordSetWriterUpdateRecord 处理器以及一个键为/canary 的属性,该属性应该将FlowFile 属性写入该键某处 在 Avro 文档中。但是,它不会出现在输出中的任何位置。

可以将 Avro 数据中的记录移动到子项中,并将元数据部分作为记录数据的一部分。不过,我不想这样做,因为它似乎不是正确的解决方案,而且它听起来比简单地修改 Avro 元数据要复杂得多。

【问题讨论】:

    标签: avro apache-nifi


    【解决方案1】:

    记录感知处理器(和读取器/写入器)不感知元数据,这意味着它们目前(从 NiFi 1.5.0 开始)不能以任何方式(检查、创建、删除等)对元数据进行操作,因此 UpdateRecord 本身不适用于元数据。使用您的 /canary 属性键,它会尝试在顶层的 Avro 记录中插入一个名为 canary 的字段,并且应该具有您指定的值。但是我相信您的输出架构需要在顶层添加金丝雀字段,否则可能会被忽略(我对此并不肯定,您可以检查输出架构以查看它是否自动添加)。

    目前没有可以显式更新 Avro 元数据的 NiFi 处理器(MergeContent 可以将各种 Avro 文件合并在一起,但您不能选择设置值,例如)。但是,我有一个未完善的 Groovy 脚本,您可以在 ExecuteScript 中使用该脚本将元数据添加到 NiFi 1.5.0+ 中的 Avro 文件。在 ExecuteScript 中,您可以将语言设置为 Groovy 并将以下语言设置为 Script Body,然后将用户定义的(也称为“动态”属性)添加到 ExecuteScript,其中键将是元数据键,并且评估值(属性支持表达式Language) 将是值:

    @Grab('org.apache.avro:avro:1.8.1')
    import org.apache.avro.*
    import org.apache.avro.file.*
    import org.apache.avro.generic.*
    
    def flowFile = session.get()
    if(!flowFile) return
    
    try {
    // Save off dynamic property values for metadata key/values later
    def metadata = [:]
    context.properties.findAll {e -> e.key.dynamic}.each {k,v -> metadata.put(k.name, context.getProperty(k).evaluateAttributeExpressions(flowFile).value.bytes)}
    
    flowFile = session.write(flowFile, {inStream, outStream ->
       DataFileStream<GenericRecord> reader = new DataFileStream<>(inStream, new GenericDatumReader<GenericRecord>())
       DataFileWriter<GenericRecord> writer = new DataFileWriter<>(new GenericDatumWriter<GenericRecord>())
       def schema = reader.schema
       def inputCodec = reader.getMetaString(DataFileConstants.CODEC) ?: DataFileConstants.NULL_CODEC
       // Forward the existing metadata to the output
       reader.metaKeys.each { key ->
          if (!DataFileWriter.isReservedMeta(key)) {
             byte[] metadatum = reader.getMeta(key)
             writer.setMeta(key, metadatum)
          }
       }
       // For each dynamic property, set the key/value pair as Avro metadata
       metadata.each {k,v -> writer.setMeta(k,v)}
       writer.setCodec(CodecFactory.fromString(inputCodec))
       writer.create(schema, outStream)
       writer.appendAllFrom(reader, false)
    } as StreamCallback)
    
    session.transfer(flowFile, REL_SUCCESS)
    } catch(e) {
       log.error('Error adding Avro metadata, penalizing flow file and routing to failure', e)
       flowFile = session.penalize(flowFile)
       session.transfer(flowFile, REL_FAILURE)
    } 
    

    请注意,此脚本可以与 1.5.0 之前的 NiFi 版本一起使用,但直到 1.5.0 才支持顶部的 @Grab,因此您必须将 Avro 及其依赖项下载到一个平面文件夹中,并在 ExecuteScript 的 Module Directory 属性中指向它。

    【讨论】:

    • 这非常有效。我很乐意看到它包含在 NiFi 中作为处理器。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-04-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多