使用 Apache Arrow C 编写 Flight 客户端:连接、上传、查询与下载完整实战
2026/9/23 7:18:25 网站建设 项目流程
  • 数据工程
  • 大数据
  • 序列化
  • 数据分析

【免费下载链接】arrow

Apache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing

项目地址:https://gitcode.com/gh_mirrors/arrow13/arrow
点击查看免费下载

本文以 Apache Arrow 仓库中 csharp/examples/FlightClientExample 示例为核心,讲解如何使用 C# 的Apache.Arrow.Flight库编写一个完整的 Arrow Flight 客户端:它能够向 Flight Asp Server 示例 上传 Arrow 表、查询元数据、下载数据并清除服务端数据。读完本文,你将掌握FlightClient各核心 API 的调用方式、RecordBatch的构造与流式传输,以及如何将客户端与服务端示例配合运行。

一、示例概述:客户端与服务端如何协作

FlightClientExample是 Apache Arrow C# 仓库(csharp/目录)中随Examples.sln提供的示例项目。它是一个命令行程序,需要配合同一个解决方案中的 FlightAspServerExample(基于 ASP.NET Core + gRPC 的 Flight 服务端)一起运行。

整个示例在 Program.cs 的Main方法中按顺序完成四件事:

  1. 上传:用StartPut把一个测试RecordBatch数组上传到服务端,以FlightDescriptor标识数据;
  2. 查询元数据:通过GetSchemaGetInfoListFlights查看已上传数据的 Schema、Flight 信息以及服务端上当前可用的所有 Flight;
  3. 下载:根据GetInfo返回的FlightInfo中的 Endpoint 与 Ticket,用GetStream把数据读回客户端;
  4. 清理:先通过ListActions发现服务端支持的动作,再用DoAction发送"clear"命令清除服务端内存中的数据。

注意:必须先启动服务端,再运行客户端。这是两个独立的 .NET 进程,客户端默认连接localhost:5000

二、项目结构、依赖与运行方式

2.1 项目文件与依赖

FlightClientExample.csproj(见 FlightClientExample.csproj)声明了:

  • OutputTypeExeTargetFrameworknet8.0
  • 通过ProjectReference直接引用仓库源码中的 Apache.Arrow.Flight.csproj,而不是 NuGet 包。

Apache.Arrow.Flight库本身依赖Google.ProtobufGrpc.Net.ClientGrpc.Tools,并引用Apache.Arrow核心库;其 gRPC 服务定义来自仓库根目录的 format/Flight.proto。因此示例使用到的FlightClientFlightDescriptorFlightInfo等类型全部位于Apache.Arrow.Flight.Client/Apache.Arrow.Flight命名空间。

2.2 编译与运行

由于两个示例都在同一解决方案下,推荐的做法是在仓库csharp/examples目录中打开 Examples.sln:

  1. 先启动FlightAspServerExample(服务端在开发环境下通过 Kestrel 监听localhost:5000的 HTTP/2 端点,详见 FlightAspServerExample/Program.cs);
  2. 再运行FlightClientExample,程序会用默认参数连接localhost:5000

客户端Main方法支持两个可选命令行参数(见 Program.cs):

FlightClientExample [host] [port]
  • host:服务端主机名,默认localhost
  • port:服务端端口,默认5000

例如连接远程服务端:

FlightClientExample 192.168.1.10 5000

三、建立连接:gRPC Channel 与 FlightClient

客户端的第一步是创建 gRPC Channel 并实例化FlightClient。Arrow Flight 本身建立在 gRPC 之上,因此客户端的一切交互都通过 gRPC Channel 完成:

// (In production systems, you should use https not http) var address = $"http://{host}:{port}"; Console.WriteLine($"Connecting to: {address}"); var channel = GrpcChannel.ForAddress(address); var client = new FlightClient(channel);

两点说明:

  • 示例出于演示目的使用http://明文连接。源码注释明确提醒:生产环境应使用 https。开发环境服务端 Program.cs 也特意配置了一个不带 TLS 的 HTTP/2 本地端点(ListenLocalhost(5000, HttpProtocols.Http2))来配合演示。
  • 从 FlightClient.cs 的构造函数可以看到,FlightClient只是对FlightService.FlightServiceClient(由Flight.proto生成的 gRPC 客户端)的一层封装,负责把Apache.Arrow.Flight的高层对象(如FlightDescriptorFlightTicket)转换为 protocol 层消息。

四、准备测试数据:用 Builder 构造多列 RecordBatch

示例通过CreateTestBatch方法构造测试数据(Program.cs),展示了RecordBatch.Builder的典型用法:

public static RecordBatch CreateTestBatch(int start, int length) { return new RecordBatch.Builder() .Append("Column A", false, col => col.Int32(array => array.AppendRange(Enumerable.Range(start, start + length)))) .Append("Column B", false, col => col.Float(array => array.AppendRange(Enumerable.Range(start, start + length).Select(x => Convert.ToSingle(x * 2))))) .Append("Column C", false, col => col.String(array => array.AppendRange(Enumerable.Range(start, start + length).Select(x => $"Item {x+1}")))) .Append("Column D", false, col => col.Boolean(array => array.AppendRange(Enumerable.Range(start, start + length).Select(x => x % 2 == 0)))) .Build(); }

该方法生成包含四列的RecordBatch

列名类型数据说明
Column AInt32start开始的连续整数,共start + length
Column BFloat整数序列的两倍并转为单精度浮点
Column CString"Item N"形式的字符串
Column DBoolean根据整数奇偶性生成布尔值

Main中构造了两个批次作为上传内容(Program.cs):

var recordBatches = new RecordBatch[] { CreateTestBatch(0, 2000), CreateTestBatch(50, 9000) };

五、上传数据:FlightDescriptor 与 StartPut

5.1 用 FlightDescriptor 标识数据

Flight 中的“数据集”由FlightDescriptor标识。它可以是名称、SQL 查询字符串或路径。示例使用命令描述符,名称为"test"(Program.cs):

var descriptor = FlightDescriptor.CreateCommandDescriptor("test");

从 FlightDescriptor.cs 源码可以看到,该类支持两种描述符:

  • CreateCommandDescriptor(string/byte[]):命令型描述符,最终映射为 protocol 层的DescriptorType.Cmd
  • CreatePathDescriptor(params string[] paths):路径型描述符,映射为DescriptorType.Path

5.2 StartPut:双工流上传

上传使用FlightClient.StartPut,它对应 gRPC 的DoPut双向流调用(FlightClient.cs):

// Upload data with StartPut var batchStreamingCall = client.StartPut(descriptor); foreach (var batch in recordBatches) { await batchStreamingCall.RequestStream.WriteAsync(batch); } // Signal we are done sending record batches await batchStreamingCall.RequestStream.CompleteAsync(); // Retrieve final response await batchStreamingCall.ResponseStream.MoveNext(); Console.WriteLine(batchStreamingCall.ResponseStream.Current.ApplicationMetadata.ToStringUtf8()); Console.WriteLine($"Wrote {recordBatches.Length} batches to server.");

关键流程:

  1. StartPut(descriptor)返回FlightRecordBatchDuplexStreamingCall,其RequestStream用于逐批写入RecordBatch
  2. 写完所有批次后必须调用RequestStream.CompleteAsync()告知服务端流结束;
  3. 服务端处理完会返回FlightPutResult(携带ApplicationMetadata),客户端通过ResponseStream.MoveNext()读取。本例服务端在收完数据后回写"Table saved."(见 InMemoryFlightServer.cs)。

从服务端实现可以看到对应的处理逻辑(InMemoryFlightServer.cs):DoPut逐条读取请求流中的批次并统计行数,随后把FlightInfo(Schema、Descriptor、Endpoint、行数)与批次列表存入FlightData内存存储。

六、查询元数据:GetSchema、GetInfo 与 ListFlights

上传完成后,客户端演示了三种元数据查询:

6.1 GetSchema —— 获取数据集的 Schema

var schema = await client.GetSchema(descriptor).ResponseAsync; Console.WriteLine($"Schema saved as: \n {schema}");

GetSchema是 unary 调用(FlightClient.cs),返回AsyncUnaryCall<Schema>,用ResponseAsync获取结果。底层会把 gRPC 响应中的 FlatBuffers Schema 字节经FlightMessageSerializer.DecodeSchema解码为Apache.Arrow.Schema。示例服务端在 InMemoryFlightServer.cs 中直接复用GetFlightInfo返回的 Schema。

6.2 GetInfo —— 获取 FlightInfo(含下载端点)

var info = await client.GetInfo(descriptor).ResponseAsync; Console.WriteLine($"Info provided: \n {info}");

GetInfo同样是一次 unary 调用(FlightClient.cs)。FlightInfo中最重要的信息是Endpoints:每个FlightEndpoint携带一个Ticket和若干候选Location,用于第 7 节的数据下载。

6.3 ListFlights —— 列出服务端所有可用 Flight

Console.WriteLine($"Available flights:"); var flights_call = client.ListFlights(); while (await flights_call.ResponseStream.MoveNext()) { Console.WriteLine(" " + flights_call.ResponseStream.Current.ToString()); }

ListFlights是服务端流式调用(FlightClient.cs),可传入FlightCriteria过滤条件(null时使用空条件)。响应流中的每一项是一个FlightInfo。示例服务端 InMemoryFlightServer.cs 会把内存中所有已注册的 Flight 逐个写出。

七、下载数据:从 FlightInfo 到 GetStream

下载是示例中最具 Flight 特色的部分。客户端并不直接从最初连接的服务器拉数据,而是依据FlightInfo.Endpoints中返回的地址与 Ticket 建立新的下载通道StreamRecordBatches方法(Program.cs)完整展示了这一模式:

public static async IAsyncEnumerable<RecordBatch> StreamRecordBatches( FlightInfo info ) { // There might be multiple endpoints hosting part of the data. In simple services, // the only endpoint might be the same server we initially queried. foreach (var endpoint in info.Endpoints) { // We may have multiple locations to choose from. Here we choose the first. var download_channel = GrpcChannel.ForAddress(endpoint.Locations.First().Uri); var download_client = new FlightClient(download_channel); var stream = download_client.GetStream(endpoint.Ticket); while (await stream.ResponseStream.MoveNext()) { yield return stream.ResponseStream.Current; } } }

要点:

  • FlightInfo.Endpoints集合:一份数据可能被切分托管在多个服务节点上,因此要遍历所有 Endpoint;
  • 每个 Endpoint 的Locations也是集合,表示同一份数据可选的多个服务地址,示例选择第一个(Locations.First().Uri);
  • GetStream(endpoint.Ticket)对应 gRPC 的DoGet(FlightClient.cs),返回FlightRecordBatchStreamingCall,通过ResponseStream.MoveNext()逐个读取RecordBatch
  • 该方法用IAsyncEnumerable<RecordBatch>暴露,因此调用方可以用await foreach消费:
await foreach (var batch in StreamRecordBatches(info)) { Console.WriteLine($"Read batch from flight server: \n {batch}"); }

配套的服务端DoGet(InMemoryFlightServer.cs)会按 Ticket 找到内存中存储的批次列表并逐个写回;若 Ticket 不存在则抛出RpcExceptionStatusCode.NotFound)。

八、发现并执行动作:ListActions 与 DoAction

Flight 协议用“动作(Action)”提供可扩展的远程命令能力。客户端先列举服务端支持的动作,再执行:

// See available commands on this server var action_stream = client.ListActions(); Console.WriteLine("Actions:"); while (await action_stream.ResponseStream.MoveNext()) { var action = action_stream.ResponseStream.Current; Console.WriteLine($" {action.Type}: {action.Description}"); } // Send clear command to drop all data from the server. var clear_result = client.DoAction(new FlightAction("clear")); await clear_result.ResponseStream.MoveNext(default);
  • ListActions():服务端流式调用(FlightClient.cs),返回FlightActionType流,包含动作类型名与描述;
  • DoAction(new FlightAction("clear")):执行动作,返回AsyncServerStreamingCall<FlightResult>(FlightClient.cs)。

示例服务端在 InMemoryFlightServer.cs 中只注册了一个"clear"动作,作用是清空内存中的所有FlightInfoRecordBatch;执行其他动作会抛出RpcExceptionStatusCode.InvalidArgument)。因此这个动作对应了本示例“收尾清理”的职责,与 readme 中描述的 4 个步骤一一对应。

九、完整运行链路小结

把客户端与服务端串起来,一次完整演示的调用链为:

FlightClientExample └─ StartPut(descriptor) → 上传 2 个 RecordBatch(DoPut) └─ GetSchema(descriptor) → 读取 Schema(GetSchema) └─ GetInfo(descriptor) → 读取 FlightInfo(GetFlightInfo) └─ ListFlights() → 列出全部 Flight └─ GetStream(info.Endpoints.Ticket) → 下载批次(DoGet,每个 Endpoint 一个连接) └─ ListActions() → 枚举动作 └─ DoAction("clear") → 清除服务端数据

每一步都能在 Program.cs 中找到对应代码,其底层 API 行为则可在 Apache.Arrow.Flight 客户端源码 中验证。这个示例虽然精简,但覆盖了 Arrow Flight 客户端最常见的全部交互模式——上传、元数据查询、下载与动作执行,是学习 Apache Arrow C# Flight 客户端 API 的最直接入口。

  • 数据工程
  • 大数据
  • 序列化
  • 数据分析

【免费下载链接】arrow

Apache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing

项目地址:https://gitcode.com/gh_mirrors/arrow13/arrow
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询