如何用arrow_array::RecordBatchReader实现类读取reqwest::Response的bytes_stream()?
实现reqwest流式读取Arrow-IPC数据并转换为RecordBatch流
依赖准备
首先在Cargo.toml中添加所需依赖,确保arrow系列 crate 版本一致:
[dependencies] reqwest = { version = "0.11", features = ["stream"] } arrow-ipc = { version = "46.0", features = ["async"] } arrow-array = "46.0" tokio = { version = "1.0", features = ["full"] } tokio-util = { version = "0.7", features = ["io"] }
核心实现代码
通过tokio-util将reqwest的字节流转换为AsyncRead,再用arrow-ipc的异步阅读器处理:
use std::sync::Arc; use arrow_array::RecordBatch; use arrow_ipc::reader::async_reader::StreamReader; use reqwest::Client; use tokio_util::io::StreamReader; #[tokio::main] async fn main() -> Result<(), Box<dyn std::error::Error>> { // 初始化reqwest客户端 let client = Client::new(); // 发送请求并获取响应 let response = client.get("https://目标API端点地址") .send() .await?; // 校验响应状态 if !response.status().is_success() { return Err(format!("请求失败,状态码: {}", response.status()).into()); } // 将reqwest的字节流转换为AsyncRead(转换错误类型以适配标准IO错误) let byte_stream = response.bytes_stream().map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e) ); let async_read = StreamReader::new(byte_stream); // 创建异步Arrow IPC流阅读器 let mut reader = StreamReader::try_new(async_read, None).await?; // 流式迭代读取RecordBatch while let Some(batch_result) = reader.next().await { let batch: Arc<RecordBatch> = batch_result?; println!("收到批次,行数: {}", batch.num_rows()); // 在这里添加你的业务处理逻辑,比如处理批次数据 // println!("批次Schema: {:?}", batch.schema()); } Ok(()) }
关键说明
- 流转换:
tokio_util::io::StreamReader是衔接reqwest异步字节流和arrow异步阅读器的核心,它把Stream<Result<Bytes>>转换为符合AsyncReadtrait的类型,同时需要将reqwest的错误转换为标准IO错误。 - 异步阅读器:
arrow-ipc::reader::async_reader::StreamReader是arrow官方提供的异步流式阅读器,专门处理Arrow-IPC格式的异步输入,支持逐批次读取RecordBatch。 - 版本兼容:确保所有arrow相关 crate(
arrow-ipc、arrow-array)使用相同版本,避免因版本差异导致的兼容性问题。
内容的提问来源于stack exchange,提问作者jwimberley
相关产品推荐
相关产品推荐

