如何在C语言中流式传输Apache Arrow RecordBatches
问题描述
- 业务场景:从PostgreSQL读取数据转换为Arrow RecordBatch后需发送至客户端,未掌握Apache Arrow C/GLib的正确使用方法。
- 测试行为:参考官方文档和源码编写最简测试样例,流程为将RecordBatch写入缓冲区、模拟缓冲区收发、从缓冲区读回RecordBatch,运行时读取失败。
- 报错现象:
- 使用写入单RecordBatch得到的buffer、或
garrow_record_batch_serialize生成的buffer构造输入流,初始化GArrowRecordBatchStreamReader时报错:[record-batch-stream-reader][open]: IOError: Expected IPC message of type schema but got record batch - 使用从缓冲区取出的裸字节构造的testBuffer创建流时,提示IPC流无效,判定数据损坏。
- 使用写入单RecordBatch得到的buffer、或
原测试代码如下:
void testRecordbatchStream(GArrowRecordBatch *rb){ GError *error = NULL; // Write Recordbatch GArrowResizableBuffer *buffer = garrow_resizable_buffer_new(300, &error); GArrowBufferOutputStream *bufferStream = garrow_buffer_output_stream_new(buffer); long written = garrow_output_stream_write_record_batch(GARROW_OUTPUT_STREAM(bufferStream), rb, NULL, &error); // Use buffer as plain bytes void *data = garrow_buffer_get_data(GARROW_BUFFER(buffer)); size_t length = garrow_buffer_get_size(GARROW_BUFFER(buffer)); // Read plain bytes and test serialize function GArrowBuffer *testBuffer = garrow_buffer_new(data, length); GArrowBuffer *arrowbuffer = garrow_record_batch_serialize(rb, NULL, &error); // Read RecordBatch from buffer GArrowBufferInputStream *inputStream = garrow_buffer_input_stream_new(arrowbuffer); GArrowRecordBatchStreamReader *sr = garrow_record_batch_stream_reader_new(GARROW_INPUT_STREAM(inputStream), &error); GArrowRecordBatch *rb2 = garrow_record_batch_reader_read_next(sr, &error); printf("Received RB: \n%s\n", garrow_record_batch_to_string(rb2, &error)); }
错误原因
核心为API使用逻辑错误,共3处:
- 格式不匹配:
garrow_output_stream_write_record_batch、garrow_record_batch_serialize两个API仅序列化单个RecordBatch的二进制内容,不会写入Arrow IPC流要求的Schema消息头。而GArrowRecordBatchStreamReader初始化时强制要求流的第一个消息为Schema,因此触发「期望schema但拿到record batch」的报错。 - 内存引用错误:直接取
garrow_buffer_get_data返回的裸指针调用garrow_buffer_new构造testBuffer时,新Buffer不会持有原内存块的引用,一旦原Buffer被释放或内存被复用,指针就会悬空,导致读取时判定IPC流损坏。 - 流未正确收尾:写入数据后没有关闭输出流/写入器,缓冲区可能存在未刷入的残留数据,导致内容不完整。
修复后实现
正确方式是使用GArrowRecordBatchStreamWriter写入完整IPC流,该组件会自动写入Schema头、RecordBatch内容、流结束标记,和StreamReader的解析逻辑完全匹配,同时注意正确处理内存引用:
void testRecordbatchStream(GArrowRecordBatch *rb){ GError *error = NULL; // 初始化可伸缩缓冲区和对应输出流 GArrowResizableBuffer *buffer = garrow_resizable_buffer_new(300, &error); GArrowBufferOutputStream *bufferStream = garrow_buffer_output_stream_new(buffer); // 提取RecordBatch的Schema,创建IPC流写入器,构造时会自动写入Schema消息头 GArrowSchema *schema = garrow_record_batch_get_schema(rb); GArrowRecordBatchStreamWriter *writer = garrow_record_batch_stream_writer_new( GARROW_OUTPUT_STREAM(bufferStream), schema, &error ); // 写入RecordBatch内容 garrow_record_batch_writer_write_record_batch( GARROW_RECORD_BATCH_WRITER(writer), rb, &error ); // 关闭写入器,自动写入流结束标记、刷新所有缓冲区内容 garrow_record_batch_writer_close( GARROW_RECORD_BATCH_WRITER(writer), &error ); g_object_unref(writer); g_object_unref(bufferStream); // 模拟网络/进程间收发:通过slice接口构造接收端Buffer,会自动持有原buffer的引用,避免悬空指针 size_t total_size = garrow_buffer_get_size(GARROW_BUFFER(buffer)); GArrowBuffer *received_buffer = garrow_buffer_slice(GARROW_BUFFER(buffer), 0, total_size); // 构造输入流、初始化流读取器 GArrowBufferInputStream *inputStream = garrow_buffer_input_stream_new(received_buffer); GArrowRecordBatchStreamReader *sr = garrow_record_batch_stream_reader_new( GARROW_INPUT_STREAM(inputStream), &error ); // 读取RecordBatch GArrowRecordBatch *rb2 = garrow_record_batch_reader_read_next(sr, &error); printf("Received RB: \n%s\n", garrow_record_batch_to_string(rb2, &error)); // 释放所有GObject引用,避免内存泄漏 g_object_unref(rb2); g_object_unref(sr); g_object_unref(inputStream); g_object_unref(received_buffer); g_object_unref(schema); g_object_unref(buffer); }
注意事项
- 不要用单RecordBatch序列化接口构造IPC流:单RecordBatch序列化结果仅包含批量数据,没有Schema和流边界标记,仅适合通信双方已经提前同步过Schema的场景下传输数据块。
- 不要直接用
garrow_buffer_get_data返回的裸指针构造新的GArrowBuffer:需要共享内存时用garrow_buffer_slice,需要独立拷贝时用garrow_buffer_copy,两者都会正确处理内存引用计数,避免悬空指针问题。 - 所有GObject实例使用完成后必须调用
g_object_unref释放,否则会出现内存泄漏。
内容的提问来源于stack exchange,提问作者Hyrikan
相关产品推荐
相关产品推荐

