【问题标题】:How to stream Apache Arrow RecordBatches in C?如何在 C 中流式传输 Apache Arrow RecordBatches?
【发布时间】:2022-07-12 17:15:28
【问题描述】:

我从 PostgreSQL 数据库中读取了一些数据,将其转换为 RecordBatches 并尝试将数据发送到客户端。但我无法正确理解 Apache Arrow C/GLib 的用法。

我的信息来源是C++ docsthe Apache Arrow C/GLib reference manualthe C/GLib Github files

通过遵循 Apache Arrow C++ 的用法描述并试验 C 中的包装类,我构建了这个将 RecordBatch 写入缓冲区并(在理论上发送和接收缓冲区之后)尝试读回该缓冲区的最小示例进入 RecordBatch。但它失败了,如果你能指出我的错误,我会很高兴!

为了便于阅读,我省略了错误捕获。代码在创建 GArrowRecordBatchStreamReader 时出错。如果我在创建 InputStream 时使用箭头缓冲区或顶部的缓冲区,则错误为 [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));
}

【问题讨论】:

    标签: c stream apache-arrow


    【解决方案1】:

    所以我的解决方案是使用 GArrowRecordBatchStreamWriter 和 Reader 类,而不是函数 garrow_output_stream_write_record_batch(),因为后者只写入没有流标头和模式的记录批。此外,必须在写入后正确访问 GArrowBuffer 的数据。 (同样,省略了错误处理)

        GError *error = NULL;
    
        GArrowResizableBuffer *buffer = garrow_resizable_buffer_new(4096, &error);
    
        GArrowBufferOutputStream *bufferStream = garrow_buffer_output_stream_new(buffer);
        GArrowSchema *schema = garrow_record_batch_get_schema(recordbatch);
        GArrowRecordBatchStreamWriter *sw = garrow_record_batch_stream_writer_new(GARROW_OUTPUT_STREAM(bufferStream), schema, &error);
    
        g_object_unref(bufferStream);
        g_object_unref(schema);
    
        gboolean test = garrow_record_batch_writer_write_record_batch(GARROW_RECORD_BATCH_WRITER(sw), recordbatch, &error);
    
        GBytes *data = garrow_buffer_get_data(GARROW_BUFFER(buffer));
        gint64 length = garrow_buffer_get_size(GARROW_BUFFER(buffer));
    
        gsize datasize;
        gconstpointer datap = g_bytes_get_data(data, &datasize);
    
        GArrowBuffer *receivingBuffer = garrow_buffer_new(datap, datasize);
    
        GArrowBufferInputStream *inputStream = garrow_buffer_input_stream_new(GARROW_BUFFER(receivingBuffer));
        GArrowRecordBatchStreamReader *sr = garrow_record_batch_stream_reader_new(GARROW_INPUT_STREAM(inputStream), &error);
    
        printf("Reading RecordBatch:\n");
        GArrowRecordBatch *recordbatch2 = garrow_record_batch_reader_read_next(GARROW_RECORD_BATCH_READER(sr), &error);
        printf("%s\n", garrow_record_batch_to_string(recordbatch2, &error));
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-11-13
      • 2020-02-04
      • 2015-12-12
      • 2014-12-31
      • 2017-11-21
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多