You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Apache NiFi中使用Protocol Buffer实现Oracle到PostgreSQL数据迁移

Apache NiFi 基于Protobuf实现Oracle到PostgreSQL数据迁移方案

核心实现逻辑

无需完全自定义整套处理器,NiFi原生支持Protobuf格式的Record读写能力,仅需替换原有Avro相关的序列化/反序列化组件即可,核心链路和原有Avro方案逻辑一致,仅替换序列化层:

原有Avro方案链路:ExecuteSQL → 输出Avro格式数据 → AvroReader + PutDatabaseRecord → PostgreSQL
替换后Protobuf方案链路:ExecuteSQL → 输出Record → ConvertRecord(AvroReader + ProtobufWriter) → 输出Protobuf格式数据 → ProtobufReader + PutDatabaseRecord → PostgreSQL

Protobuf为强schema约束协议,需要提前编写和迁移表结构匹配的.proto schema文件,这是实现的前置核心依赖。

前置准备

  • 为每一张需要迁移的Oracle表编写对应Protobuf schema(.proto文件),字段类型需和Oracle、PostgreSQL的字段类型做好映射,比如Oracle的NUMBER对应Protobuf的int64/double,VARCHAR2对应string等
  • 使用protoc编译器把.proto文件编译为Java类,后续NiFi的Protobuf读写组件需要依赖该类完成序列化反序列化

操作步骤

  1. 配置Oracle侧数据提取处理器
    • 保留原有ExecuteSQL处理器,数据库连接池配置保持不变,将「输出格式」调整为Record
  2. 配置Protobuf序列化组件
    • 添加ConvertRecord处理器,和ExecuteSQL建立下游连接
    • 新增AvroReader服务,用于读取ExecuteSQL输出的原生Record元数据
    • 新增ProtobufWriter服务,核心配置项:
      • Protobuf Message Type:填写提前编译好的Protobuf Java类全路径
      • Schema Access Strategy:选择Use Embedded Schema或者关联上传的.proto文件
  3. (可选)如果需要中间落盘或者跨节点传输Protobuf格式数据,正常使用PutFile/FetchFile或者队列传输即可,此时数据已经是Protobuf序列化后的二进制格式
  4. 配置入库处理器
    • 原有PutDatabaseRecord处理器的Reader服务替换为新配置的ProtobufReader,ProtobufReader配置和ProtobufWriter保持一致,关联相同的.proto schema和Message类
    • PostgreSQL连接池、批量提交等其余配置保持不变即可

自定义处理器伪代码参考(适用于定制化扩展场景)

如果需要自行实现更灵活的Protobuf处理逻辑,可参考如下伪代码:

import com.google.protobuf.Descriptors.Descriptor;
import com.google.protobuf.Descriptors.FieldDescriptor;
import com.google.protobuf.DynamicMessage;
import org.apache.nifi.processor.*;
import org.apache.nifi.record.Record;
import org.apache.nifi.record.RecordReader;
import org.apache.nifi.record.RecordReaderFactory;

// 处理器核心属性定义
public static final PropertyDescriptor PROTO_MESSAGE_CLASS = new PropertyDescriptor.Builder()
    .name("Protobuf Message全类名")
    .required(true)
    .build();
public static final PropertyDescriptor RECORD_READER = new PropertyDescriptor.Builder()
    .name("上游Record读取服务")
    .identifiesControllerService(RecordReaderFactory.class)
    .required(true)
    .build();
public static final Relationship REL_SUCCESS = new Relationship.Builder().name("success").build();
public static final Relationship REL_FAILURE = new Relationship.Builder().name("failure").build();

private Descriptor protoDescriptor;
private DynamicMessage.Builder protobufMessageBuilder;

// 处理器初始化时加载Protobuf schema
@OnScheduled
public void init(ProcessContext context) throws Exception {
    Class<?> protoClass = Class.forName(context.getProperty(PROTO_MESSAGE_CLASS).getValue());
    Method getDescriptorMethod = protoClass.getMethod("getDescriptor");
    protoDescriptor = (Descriptor) getDescriptorMethod.invoke(null);
    protobufMessageBuilder = DynamicMessage.newBuilder(protoDescriptor);
}

@Override
public void onTrigger(ProcessContext context, ProcessSession session) throws ProcessException {
    FlowFile flowFile = session.get();
    if (flowFile == null) return;

    RecordReaderFactory readerFactory = context.getProperty(RECORD_READER).asControllerService(RecordReaderFactory.class);
    try (InputStream in = session.read(flowFile);
         RecordReader reader = readerFactory.createRecordReader(flowFile, in, flowFile.getSize(), getLogger())) {
        ByteArrayOutputStream protoOut = new ByteArrayOutputStream();
        Record record;
        // 逐个读取上游数据库Record,转换为Protobuf格式
        while ((record = reader.nextRecord()) != null) {
            for (FieldDescriptor fieldDesc : protoDescriptor.getFields()) {
                Object fieldValue = record.getValue(fieldDesc.getName());
                // 按需添加字段类型转换逻辑
                if (fieldDesc.getType() == FieldDescriptor.Type.INT64 && fieldValue instanceof Number) {
                    fieldValue = ((Number) fieldValue).longValue();
                }
                protobufMessageBuilder.setField(fieldDesc, fieldValue);
            }
            // 序列化Protobuf消息写入输出流
            protobufMessageBuilder.build().writeDelimitedTo(protoOut);
            protobufMessageBuilder.clear();
        }
        // 生成Protobuf格式的FlowFile传递到下游
        FlowFile protoFlowFile = session.create(flowFile);
        protoFlowFile = session.write(protoFlowFile, out -> out.write(protoOut.toByteArray()));
        session.transfer(protoFlowFile, REL_SUCCESS);
        session.remove(flowFile);
    } catch (Exception e) {
        getLogger().error("Protobuf序列化处理失败", e);
        session.transfer(flowFile, REL_FAILURE);
    }
}

注意事项

  • Protobuf字段顺序、类型必须和.proto文件完全一致,否则会出现反序列化失败问题
  • 如果表结构发生变更,需要同步更新.proto文件并重新编译Java类,同时更新NiFi中Protobuf读写服务的配置
  • 批量迁移场景下可适当调大ProtobufWriter的缓冲区大小,提升序列化性能

内容的提问来源于stack exchange,提问作者vigneshwar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.27 09:45:08