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

如何在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流无效,判定数据损坏。

原测试代码如下:

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处:

  1. 格式不匹配:garrow_output_stream_write_record_batch、garrow_record_batch_serialize两个API仅序列化单个RecordBatch的二进制内容,不会写入Arrow IPC流要求的Schema消息头。而GArrowRecordBatchStreamReader初始化时强制要求流的第一个消息为Schema,因此触发「期望schema但拿到record batch」的报错。
  2. 内存引用错误:直接取garrow_buffer_get_data返回的裸指针调用garrow_buffer_new构造testBuffer时,新Buffer不会持有原内存块的引用,一旦原Buffer被释放或内存被复用,指针就会悬空,导致读取时判定IPC流损坏。
  3. 流未正确收尾:写入数据后没有关闭输出流/写入器,缓冲区可能存在未刷入的残留数据,导致内容不完整。
修复后实现

正确方式是使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 11:36:19