如何在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读写组件需要依赖该类完成序列化反序列化
操作步骤
- 配置Oracle侧数据提取处理器
- 保留原有
ExecuteSQL处理器,数据库连接池配置保持不变,将「输出格式」调整为Record
- 保留原有
- 配置Protobuf序列化组件
- 添加
ConvertRecord处理器,和ExecuteSQL建立下游连接 - 新增
AvroReader服务,用于读取ExecuteSQL输出的原生Record元数据 - 新增
ProtobufWriter服务,核心配置项:Protobuf Message Type:填写提前编译好的Protobuf Java类全路径Schema Access Strategy:选择Use Embedded Schema或者关联上传的.proto文件
- 添加
- (可选)如果需要中间落盘或者跨节点传输Protobuf格式数据,正常使用
PutFile/FetchFile或者队列传输即可,此时数据已经是Protobuf序列化后的二进制格式 - 配置入库处理器
- 原有
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
相关产品推荐
相关产品推荐

