记录感知处理器(和读取器/写入器)不感知元数据,这意味着它们目前(从 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 属性中指向它。