- 数据工程
- 大数据
- 序列化
- 数据分析
【免费下载链接】arrow
Apache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing
本文以 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方法中按顺序完成四件事:
- 上传:用
StartPut把一个测试RecordBatch数组上传到服务端,以FlightDescriptor标识数据; - 查询元数据:通过
GetSchema、GetInfo和ListFlights查看已上传数据的 Schema、Flight 信息以及服务端上当前可用的所有 Flight; - 下载:根据
GetInfo返回的FlightInfo中的 Endpoint 与 Ticket,用GetStream把数据读回客户端; - 清理:先通过
ListActions发现服务端支持的动作,再用DoAction发送"clear"命令清除服务端内存中的数据。
注意:必须先启动服务端,再运行客户端。这是两个独立的 .NET 进程,客户端默认连接
localhost:5000。
二、项目结构、依赖与运行方式
2.1 项目文件与依赖
FlightClientExample.csproj(见 FlightClientExample.csproj)声明了:
OutputType为Exe,TargetFramework为net8.0;- 通过
ProjectReference直接引用仓库源码中的 Apache.Arrow.Flight.csproj,而不是 NuGet 包。
Apache.Arrow.Flight库本身依赖Google.Protobuf、Grpc.Net.Client与Grpc.Tools,并引用Apache.Arrow核心库;其 gRPC 服务定义来自仓库根目录的 format/Flight.proto。因此示例使用到的FlightClient、FlightDescriptor、FlightInfo等类型全部位于Apache.Arrow.Flight.Client/Apache.Arrow.Flight命名空间。
2.2 编译与运行
由于两个示例都在同一解决方案下,推荐的做法是在仓库csharp/examples目录中打开 Examples.sln:
- 先启动
FlightAspServerExample(服务端在开发环境下通过 Kestrel 监听localhost:5000的 HTTP/2 端点,详见 FlightAspServerExample/Program.cs); - 再运行
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的高层对象(如FlightDescriptor、FlightTicket)转换为 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 A | Int32 | 从start开始的连续整数,共start + length个 |
| Column B | Float | 整数序列的两倍并转为单精度浮点 |
| Column C | String | "Item N"形式的字符串 |
| Column D | Boolean | 根据整数奇偶性生成布尔值 |
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.");关键流程:
StartPut(descriptor)返回FlightRecordBatchDuplexStreamingCall,其RequestStream用于逐批写入RecordBatch;- 写完所有批次后必须调用
RequestStream.CompleteAsync()告知服务端流结束; - 服务端处理完会返回
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 不存在则抛出RpcException(StatusCode.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"动作,作用是清空内存中的所有FlightInfo与RecordBatch;执行其他动作会抛出RpcException(StatusCode.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
相关推荐
Apache Arrow .NET Flight 客户端实战:基于 Apache.Arrow.Flight 实现表数据的上传、查询与下载
Apache Arrow .NET Flight 客户端实战:基于 Apache.Arrow.Flight 实现表数据的上传、查询与下载 导读 本文以 Apac
数据工程数据分析大数据Apache Arrow C++ 中编写 Flight RPC 服务:服务端、客户端、认证与最佳实践
Apache Arrow C++ 中编写 Flight RPC 服务:服务端、客户端、认证与最佳实践 导读 :Arrow Flight 是 Apache Arr
数据工程数据分析大数据Apache Arrow Flight Ruby 绑定:Red Arrow Flight 的 GObject Introspection 桥接原理与客户端实战
Apache Arrow Flight Ruby 绑定:Red Arrow Flight 的 GObject Introspection 桥接原理与客户端实战
大数据数据分析数据工程序列化
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考