数据工程大数据序列化数据分析【免费下载链接】arrowApache 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方法中按顺序完成四件事上传用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为ExeTargetFramework为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.csFlightClientExample [host] [port]host服务端主机名默认localhostport服务端端口默认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 {x1})))) .Append(Column D, false, col col.Boolean(array array.AppendRange(Enumerable.Range(start, start length).Select(x x % 2 0)))) .Build(); }该方法生成包含四列的RecordBatch列名类型数据说明Column AInt32从start开始的连续整数共start length个Column BFloat整数序列的两倍并转为单精度浮点Column CStringItem N形式的字符串Column DBoolean根据整数奇偶性生成布尔值Main中构造了两个批次作为上传内容Program.csvar recordBatches new RecordBatch[] { CreateTestBatch(0, 2000), CreateTestBatch(50, 9000) };五、上传数据FlightDescriptor 与 StartPut5.1 用 FlightDescriptor 标识数据Flight 中的“数据集”由FlightDescriptor标识。它可以是名称、SQL 查询字符串或路径。示例使用命令描述符名称为testProgram.csvar descriptor FlightDescriptor.CreateCommandDescriptor(test);从 FlightDescriptor.cs 源码可以看到该类支持两种描述符CreateCommandDescriptor(string/byte[])命令型描述符最终映射为 protocol 层的DescriptorType.CmdCreatePathDescriptor(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.csDoPut逐条读取请求流中的批次并统计行数随后把FlightInfoSchema、Descriptor、Endpoint、行数与批次列表存入FlightData内存存储。六、查询元数据GetSchema、GetInfo 与 ListFlights上传完成后客户端演示了三种元数据查询6.1 GetSchema —— 获取数据集的 Schemavar schema await client.GetSchema(descriptor).ResponseAsync; Console.WriteLine($Schema saved as: \n {schema});GetSchema是 unary 调用FlightClient.cs返回AsyncUnaryCallSchema用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 —— 列出服务端所有可用 FlightConsole.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 IAsyncEnumerableRecordBatch 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().UriGetStream(endpoint.Ticket)对应 gRPC 的DoGetFlightClient.cs返回FlightRecordBatchStreamingCall通过ResponseStream.MoveNext()逐个读取RecordBatch该方法用IAsyncEnumerableRecordBatch暴露因此调用方可以用await foreach消费await foreach (var batch in StreamRecordBatches(info)) { Console.WriteLine($Read batch from flight server: \n {batch}); }配套的服务端DoGetInMemoryFlightServer.cs会按 Ticket 找到内存中存储的批次列表并逐个写回若 Ticket 不存在则抛出RpcExceptionStatusCode.NotFound。八、发现并执行动作ListActions 与 DoActionFlight 协议用“动作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))执行动作返回AsyncServerStreamingCallFlightResultFlightClient.cs。示例服务端在 InMemoryFlightServer.cs 中只注册了一个clear动作作用是清空内存中的所有FlightInfo与RecordBatch执行其他动作会抛出RpcExceptionStatusCode.InvalidArgument。因此这个动作对应了本示例“收尾清理”的职责与 readme 中描述的 4 个步骤一一对应。九、完整运行链路小结把客户端与服务端串起来一次完整演示的调用链为FlightClientExample └─ StartPut(descriptor) → 上传 2 个 RecordBatchDoPut └─ GetSchema(descriptor) → 读取 SchemaGetSchema └─ GetInfo(descriptor) → 读取 FlightInfoGetFlightInfo └─ ListFlights() → 列出全部 Flight └─ GetStream(info.Endpoints.Ticket) → 下载批次DoGet每个 Endpoint 一个连接 └─ ListActions() → 枚举动作 └─ DoAction(clear) → 清除服务端数据每一步都能在 Program.cs 中找到对应代码其底层 API 行为则可在 Apache.Arrow.Flight 客户端源码 中验证。这个示例虽然精简但覆盖了 Arrow Flight 客户端最常见的全部交互模式——上传、元数据查询、下载与动作执行是学习 Apache Arrow C# Flight 客户端 API 的最直接入口。赞分享数据工程大数据序列化数据分析【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址https://gitcode.com/gh_mirrors/arrow13/arrow点击查看免费下载相关推荐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),仅供参考
