From 02b1164643e8111ac586304c6f77a08e700e0960 Mon Sep 17 00:00:00 2001 From: Faiaz Sanaulla Date: Thu, 8 May 2025 12:39:13 +0200 Subject: [PATCH 1/5] add more logs to flight decoder/encoder --- arrow-flight/Cargo.toml | 1 + arrow-flight/examples/flight_sql_server.rs | 10 ++-- arrow-flight/src/client.rs | 6 +- arrow-flight/src/decode.rs | 18 ++++-- arrow-flight/src/encode.rs | 56 ++++++++++++------ arrow-flight/src/sql/client.rs | 10 ++-- arrow-flight/tests/client.rs | 6 +- arrow-flight/tests/common/server.rs | 2 +- arrow-flight/tests/encode_decode.rs | 64 +++++++++++---------- arrow-flight/tests/flight_sql_client.rs | 1 + arrow-flight/tests/flight_sql_client_cli.rs | 9 +-- 11 files changed, 114 insertions(+), 69 deletions(-) diff --git a/arrow-flight/Cargo.toml b/arrow-flight/Cargo.toml index 9531c63af003..45880d644c6f 100644 --- a/arrow-flight/Cargo.toml +++ b/arrow-flight/Cargo.toml @@ -49,6 +49,7 @@ prost = { version = "0.13.1", default-features = false, features = ["prost-deriv prost-types = { version = "0.13.1", default-features = false } tokio = { version = "1.0", default-features = false, features = ["macros", "rt", "rt-multi-thread"] } tonic = { version = "0.13", default-features = false, features = ["transport", "codegen", "prost"] } +tracing = { version = "0.1.40", features = ["default", "log"] } # CLI-related dependencies anyhow = { version = "1.0", optional = true } diff --git a/arrow-flight/examples/flight_sql_server.rs b/arrow-flight/examples/flight_sql_server.rs index 657298b4a8b3..7fccb11c0a23 100644 --- a/arrow-flight/examples/flight_sql_server.rs +++ b/arrow-flight/examples/flight_sql_server.rs @@ -458,7 +458,7 @@ impl FlightSqlService for FlightSqlServiceImpl { let batch = builder.build(); let stream = FlightDataEncoderBuilder::new() .with_schema(schema) - .build(futures::stream::once(async { batch })) + .build(futures::stream::once(async { batch }), "") .map_err(Status::from); Ok(Response::new(Box::pin(stream))) } @@ -485,7 +485,7 @@ impl FlightSqlService for FlightSqlServiceImpl { let batch = builder.build(); let stream = FlightDataEncoderBuilder::new() .with_schema(schema) - .build(futures::stream::once(async { batch })) + .build(futures::stream::once(async { batch }), "") .map_err(Status::from); Ok(Response::new(Box::pin(stream))) } @@ -525,7 +525,7 @@ impl FlightSqlService for FlightSqlServiceImpl { let batch = builder.build(); let stream = FlightDataEncoderBuilder::new() .with_schema(schema) - .build(futures::stream::once(async { batch })) + .build(futures::stream::once(async { batch }), "") .map_err(Status::from); Ok(Response::new(Box::pin(stream))) } @@ -548,7 +548,7 @@ impl FlightSqlService for FlightSqlServiceImpl { let batch = builder.build(); let stream = FlightDataEncoderBuilder::new() .with_schema(schema) - .build(futures::stream::once(async { batch })) + .build(futures::stream::once(async { batch }), "") .map_err(Status::from); Ok(Response::new(Box::pin(stream))) } @@ -602,7 +602,7 @@ impl FlightSqlService for FlightSqlServiceImpl { let batch = builder.build(); let stream = FlightDataEncoderBuilder::new() .with_schema(schema) - .build(futures::stream::once(async { batch })) + .build(futures::stream::once(async { batch }), "") .map_err(Status::from); Ok(Response::new(Box::pin(stream))) } diff --git a/arrow-flight/src/client.rs b/arrow-flight/src/client.rs index 97d9899a9fb0..5b4b111530c4 100644 --- a/arrow-flight/src/client.rs +++ b/arrow-flight/src/client.rs @@ -211,6 +211,7 @@ impl FlightClient { Ok(FlightRecordBatchStream::new_from_flight_data( response_stream.map_err(FlightError::Tonic), + "", ) .with_headers(md) .with_trailers(trailers)) @@ -429,7 +430,10 @@ impl FlightClient { let error_stream = FallibleTonicResponseStream::new(receiver, response_stream); // combine the response from the server and any error from the client - Ok(FlightRecordBatchStream::new_from_flight_data(error_stream)) + Ok(FlightRecordBatchStream::new_from_flight_data( + error_stream, + "", + )) } /// Make a `ListFlights` call to the server with the provided diff --git a/arrow-flight/src/decode.rs b/arrow-flight/src/decode.rs index 7bafc384306b..7dae90d000ca 100644 --- a/arrow-flight/src/decode.rs +++ b/arrow-flight/src/decode.rs @@ -23,6 +23,7 @@ use bytes::Bytes; use futures::{ready, stream::BoxStream, Stream, StreamExt}; use std::{collections::HashMap, fmt::Debug, pin::Pin, sync::Arc, task::Poll}; use tonic::metadata::MetadataMap; +use tracing::debug; use crate::error::{FlightError, Result}; @@ -101,12 +102,12 @@ impl FlightRecordBatchStream { } /// Create a new [`FlightRecordBatchStream`] from a stream of [`FlightData`] - pub fn new_from_flight_data(inner: S) -> Self + pub fn new_from_flight_data(inner: S, reader_id: &str) -> Self where S: Stream> + Send + 'static, { Self { - inner: FlightDataDecoder::new(inner), + inner: FlightDataDecoder::new(inner, reader_id), headers: MetadataMap::default(), trailers: None, } @@ -234,6 +235,8 @@ pub struct FlightDataDecoder { state: Option, /// Seen the end of the inner stream? done: bool, + + reader_id: String, } impl Debug for FlightDataDecoder { @@ -248,7 +251,7 @@ impl Debug for FlightDataDecoder { impl FlightDataDecoder { /// Create a new wrapper around the stream of [`FlightData`] - pub fn new(response: S) -> Self + pub fn new(response: S, reader_id: &str) -> Self where S: Stream> + Send + 'static, { @@ -256,6 +259,7 @@ impl FlightDataDecoder { state: None, response: response.boxed(), done: false, + reader_id: reader_id.to_string(), } } @@ -354,20 +358,26 @@ impl futures::Stream for FlightDataDecoder { cx: &mut std::task::Context<'_>, ) -> Poll> { if self.done { + debug!(self.reader_id, "stream done"); return Poll::Ready(None); } loop { + debug!(self.reader_id, "polling next message"); let res = ready!(self.response.poll_next_unpin(cx)); return Poll::Ready(match res { None => { self.done = true; + debug!(self.reader_id, "inner is exhausted"); None // inner is exhausted } Some(data) => Some(match data { Err(e) => Err(e), Ok(data) => match self.extract_message(data) { - Ok(Some(extracted)) => Ok(extracted), + Ok(Some(extracted)) => { + debug!(self.reader_id, "message extracted"); + Ok(extracted) + } Ok(None) => continue, // Need next input message Err(e) => Err(e), }, diff --git a/arrow-flight/src/encode.rs b/arrow-flight/src/encode.rs index 561c26fa8282..b7ba68e1f870 100644 --- a/arrow-flight/src/encode.rs +++ b/arrow-flight/src/encode.rs @@ -25,6 +25,7 @@ use arrow_ipc::writer::{DictionaryTracker, IpcDataGenerator, IpcWriteOptions}; use arrow_schema::{DataType, Field, FieldRef, Fields, Schema, SchemaRef, UnionMode}; use bytes::Bytes; use futures::{ready, stream::BoxStream, Stream, StreamExt}; +use tracing::debug; /// Creates a [`Stream`] of [`FlightData`]s from a /// `Stream` of [`Result`]<[`RecordBatch`], [`FlightError`]>. @@ -238,7 +239,7 @@ impl FlightDataEncoderBuilder { /// of [`FlightData`], consuming self. /// /// See example on [`Self`] and [`FlightDataEncoder`] for more details - pub fn build(self, input: S) -> FlightDataEncoder + pub fn build(self, input: S, reader_id: &str) -> FlightDataEncoder where S: Stream> + Send + 'static, { @@ -259,6 +260,7 @@ impl FlightDataEncoderBuilder { app_metadata, descriptor, dictionary_handling, + reader_id, ) } } @@ -287,9 +289,12 @@ pub struct FlightDataEncoder { /// Deterimines how `DictionaryArray`s are encoded for transport. /// See [`DictionaryHandling`] for more information. dictionary_handling: DictionaryHandling, + + reader_id: String, } impl FlightDataEncoder { + #[allow(clippy::too_many_arguments)] fn new( inner: BoxStream<'static, Result>, schema: Option, @@ -298,6 +303,7 @@ impl FlightDataEncoder { app_metadata: Bytes, descriptor: Option, dictionary_handling: DictionaryHandling, + reader_id: &str, ) -> Self { let mut encoder = Self { inner, @@ -312,6 +318,7 @@ impl FlightDataEncoder { done: false, descriptor, dictionary_handling, + reader_id: reader_id.to_string(), }; // If schema is known up front, enqueue it immediately @@ -398,12 +405,21 @@ impl Stream for FlightDataEncoder { cx: &mut std::task::Context<'_>, ) -> Poll> { loop { + debug!( + reader_id = self.reader_id, + "flight data encoder polling next" + ); if self.done && self.queue.is_empty() { + debug!( + reader_id = self.reader_id, + "stream done, no more data to send" + ); return Poll::Ready(None); } // Any messages queued to send? if let Some(data) = self.queue.pop_front() { + debug!(reader_id = self.reader_id, "sending queued message"); return Poll::Ready(Some(Ok(data))); } @@ -416,6 +432,10 @@ impl Stream for FlightDataEncoder { self.done = true; // queue must also be empty so we are done assert!(self.queue.is_empty()); + debug!( + reader_id = self.reader_id, + "stream done, no more data to send" + ); return Poll::Ready(None); } Some(Err(e)) => { @@ -425,6 +445,7 @@ impl Stream for FlightDataEncoder { return Poll::Ready(Some(Err(e))); } Some(Ok(batch)) => { + debug!(reader_id = self.reader_id, "got batch"); // had data, encode into the queue if let Err(e) = self.encode_batch(batch) { self.done = true; @@ -790,8 +811,8 @@ mod tests { let stream = futures::stream::iter(vec![Ok(batch1), Ok(batch2)]); - let encoder = FlightDataEncoderBuilder::default().build(stream); - let mut decoder = FlightDataDecoder::new(encoder); + let encoder = FlightDataEncoderBuilder::default().build(stream, ""); + let mut decoder = FlightDataDecoder::new(encoder, ""); let expected_schema = Schema::new(vec![Field::new("dict", DataType::Utf8, false)]); let expected_schema = Arc::new(expected_schema); let mut expected_arrays = vec![ @@ -851,7 +872,7 @@ mod tests { let encoder = FlightDataEncoderBuilder::default() .with_schema(schema) - .build(stream); + .build(stream, ""); let expected_schema = Arc::new(Schema::new(vec![Field::new("dict", DataType::Utf8, false)])); assert_eq!(Some(expected_schema), encoder.known_schema()) @@ -876,7 +897,7 @@ mod tests { let encoder = FlightDataEncoderBuilder::default() .with_dictionary_handling(DictionaryHandling::Resend) .with_schema(schema.clone()) - .build(stream); + .build(stream, ""); assert_eq!(Some(schema), encoder.known_schema()) } @@ -929,9 +950,9 @@ mod tests { let stream = futures::stream::iter(vec![Ok(batch1), Ok(batch2)]); - let encoder = FlightDataEncoderBuilder::default().build(stream); + let encoder = FlightDataEncoderBuilder::default().build(stream, ""); - let mut decoder = FlightDataDecoder::new(encoder); + let mut decoder = FlightDataDecoder::new(encoder, ""); let expected_schema = Schema::new(vec![Field::new_list( "dict_list", Field::new("item", DataType::Utf8, true), @@ -1031,9 +1052,9 @@ mod tests { let stream = futures::stream::iter(vec![Ok(batch1), Ok(batch2)]); - let encoder = FlightDataEncoderBuilder::default().build(stream); + let encoder = FlightDataEncoderBuilder::default().build(stream, ""); - let mut decoder = FlightDataDecoder::new(encoder); + let mut decoder = FlightDataDecoder::new(encoder, ""); let expected_schema = Schema::new(vec![Field::new_struct( "struct", vec![Field::new_list( @@ -1212,9 +1233,9 @@ mod tests { let stream = futures::stream::iter(vec![Ok(batch1), Ok(batch2), Ok(batch3)]); - let encoder = FlightDataEncoderBuilder::default().build(stream); + let encoder = FlightDataEncoderBuilder::default().build(stream, ""); - let mut decoder = FlightDataDecoder::new(encoder); + let mut decoder = FlightDataDecoder::new(encoder, ""); let hydrated_struct_fields = vec![Field::new_list( "dict_list", @@ -1427,9 +1448,9 @@ mod tests { let stream = futures::stream::iter(vec![Ok(batch1), Ok(batch2)]); - let encoder = FlightDataEncoderBuilder::default().build(stream); + let encoder = FlightDataEncoderBuilder::default().build(stream, ""); - let mut decoder = FlightDataDecoder::new(encoder); + let mut decoder = FlightDataDecoder::new(encoder, ""); let expected_schema = Schema::new(vec![Field::new_map( "dict_map", "entries", @@ -1540,11 +1561,14 @@ mod tests { let encoder = FlightDataEncoderBuilder::default() .with_options(IpcWriteOptions::default().with_preserve_dict_id(false)) .with_dictionary_handling(DictionaryHandling::Resend) - .build(futures::stream::iter(batches.clone().into_iter().map(Ok))); + .build( + futures::stream::iter(batches.clone().into_iter().map(Ok)), + "", + ); let mut expected_batches = batches.drain(..); - let mut decoder = FlightDataDecoder::new(encoder); + let mut decoder = FlightDataDecoder::new(encoder, ""); while let Some(decoded) = decoder.next().await { let decoded = decoded.unwrap(); match decoded.payload { @@ -1841,7 +1865,7 @@ mod tests { .with_max_flight_data_size(max_flight_data_size) // use 8-byte alignment - default alignment is 64 which produces bigger ipc data .with_options(IpcWriteOptions::try_new(8, false, MetadataVersion::V5).unwrap()) - .build(futures::stream::iter([Ok(batch.clone())])); + .build(futures::stream::iter([Ok(batch.clone())]), ""); let mut i = 0; while let Some(data) = stream.next().await.transpose().unwrap() { diff --git a/arrow-flight/src/sql/client.rs b/arrow-flight/src/sql/client.rs index e45e505b2b61..df6e466a28d6 100644 --- a/arrow-flight/src/sql/client.rs +++ b/arrow-flight/src/sql/client.rs @@ -247,7 +247,7 @@ impl FlightSqlServiceClient { let descriptor = FlightDescriptor::new_cmd(command.as_any().encode_to_vec()); let flight_data = FlightDataEncoderBuilder::new() .with_flight_descriptor(Some(descriptor)) - .build(stream); + .build(stream, ""); // Intercept client errors and send them to the one shot channel above let flight_data = Box::pin(flight_data); @@ -310,6 +310,7 @@ impl FlightSqlServiceClient { Ok(FlightRecordBatchStream::new_from_flight_data( response_stream.map_err(FlightError::Tonic), + "", ) .with_headers(md) .with_trailers(trailers)) @@ -627,9 +628,10 @@ impl PreparedStatement { .with_flight_descriptor(Some(descriptor)) .with_schema(params_batch.schema()); let flight_data = flight_stream_builder - .build(futures::stream::iter( - self.parameter_binding.clone().map(Ok), - )) + .build( + futures::stream::iter(self.parameter_binding.clone().map(Ok)), + "", + ) .try_collect::>() .await .map_err(flight_error_to_arrow_error)?; diff --git a/arrow-flight/tests/client.rs b/arrow-flight/tests/client.rs index 25dad0e77a3e..fa10ded13017 100644 --- a/arrow-flight/tests/client.rs +++ b/arrow-flight/tests/client.rs @@ -498,7 +498,7 @@ async fn test_do_exchange() { let expected_stream = futures::stream::iter(output_flight_data).map(Ok); let expected_batches: Vec<_> = - FlightRecordBatchStream::new_from_flight_data(expected_stream) + FlightRecordBatchStream::new_from_flight_data(expected_stream, "") .try_collect() .await .unwrap(); @@ -1088,7 +1088,7 @@ async fn test_flight_data() -> Vec { // encode the batch as a stream of FlightData FlightDataEncoderBuilder::new() - .build(futures::stream::iter(vec![Ok(batch)])) + .build(futures::stream::iter(vec![Ok(batch)]), "") .try_collect() .await .unwrap() @@ -1103,7 +1103,7 @@ async fn test_flight_data2() -> Vec { // encode the batch as a stream of FlightData FlightDataEncoderBuilder::new() - .build(futures::stream::iter(vec![Ok(batch)])) + .build(futures::stream::iter(vec![Ok(batch)]), "") .try_collect() .await .unwrap() diff --git a/arrow-flight/tests/common/server.rs b/arrow-flight/tests/common/server.rs index a004ccb0737e..b9f29fb462b7 100644 --- a/arrow-flight/tests/common/server.rs +++ b/arrow-flight/tests/common/server.rs @@ -410,7 +410,7 @@ impl FlightService for TestFlightServer { let batch_stream = futures::stream::iter(batches).map_err(Into::into); let stream = FlightDataEncoderBuilder::new() - .build(batch_stream) + .build(batch_stream, "") .map_err(Into::into); let mut resp = Response::new(stream.boxed()); diff --git a/arrow-flight/tests/encode_decode.rs b/arrow-flight/tests/encode_decode.rs index cbfae1825845..608455562a07 100644 --- a/arrow-flight/tests/encode_decode.rs +++ b/arrow-flight/tests/encode_decode.rs @@ -53,9 +53,9 @@ async fn test_error() { futures::stream::iter(vec![Err(FlightError::NotYetImplemented("foo".into()))]); let encoder = FlightDataEncoderBuilder::default(); - let encode_stream = encoder.build(input_batch_stream); + let encode_stream = encoder.build(input_batch_stream, ""); - let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream); + let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream, ""); let result: Result, _> = decode_stream.try_collect().await; let result = result.unwrap_err(); @@ -131,9 +131,9 @@ async fn test_view_types_many() { #[tokio::test] async fn test_zero_batches_no_schema() { - let stream = FlightDataEncoderBuilder::default().build(futures::stream::iter(vec![])); + let stream = FlightDataEncoderBuilder::default().build(futures::stream::iter(vec![]), ""); - let mut decoder = FlightRecordBatchStream::new_from_flight_data(stream); + let mut decoder = FlightRecordBatchStream::new_from_flight_data(stream, ""); assert!(decoder.schema().is_none()); // No batches come out assert!(decoder.next().await.is_none()); @@ -146,9 +146,9 @@ async fn test_zero_batches_schema_specified() { let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int64, false)])); let stream = FlightDataEncoderBuilder::default() .with_schema(schema.clone()) - .build(futures::stream::iter(vec![])); + .build(futures::stream::iter(vec![]), ""); - let mut decoder = FlightRecordBatchStream::new_from_flight_data(stream); + let mut decoder = FlightRecordBatchStream::new_from_flight_data(stream, ""); assert!(decoder.schema().is_none()); // No batches come out assert!(decoder.next().await.is_none()); @@ -171,7 +171,7 @@ async fn test_with_flight_descriptor() { .with_schema(schema.clone()) .with_flight_descriptor(descriptor.clone()); - let mut encoder = encoder.build(stream); + let mut encoder = encoder.build(stream, ""); // First batch should be the schema let first_batch = encoder.next().await.unwrap().unwrap(); @@ -193,9 +193,9 @@ async fn test_zero_batches_dictionary_schema_specified() { ])); let stream = FlightDataEncoderBuilder::default() .with_schema(schema.clone()) - .build(futures::stream::iter(vec![])); + .build(futures::stream::iter(vec![]), ""); - let mut decoder = FlightRecordBatchStream::new_from_flight_data(stream); + let mut decoder = FlightRecordBatchStream::new_from_flight_data(stream, ""); assert!(decoder.schema().is_none()); // No batches come out assert!(decoder.next().await.is_none()); @@ -210,10 +210,11 @@ async fn test_app_metadata() { let app_metadata = Bytes::from("My Metadata"); let encoder = FlightDataEncoderBuilder::default().with_metadata(app_metadata.clone()); - let encode_stream = encoder.build(input_batch_stream); + let encode_stream = encoder.build(input_batch_stream, ""); // use lower level stream to get access to app metadata - let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream).into_inner(); + let decode_stream = + FlightRecordBatchStream::new_from_flight_data(encode_stream, "").into_inner(); let mut messages: Vec<_> = decode_stream.try_collect().await.expect("encode fails"); @@ -239,10 +240,11 @@ async fn test_max_message_size() { // 5 input rows, with a very small limit should result in 5 batch messages let encoder = FlightDataEncoderBuilder::default().with_max_flight_data_size(1); - let encode_stream = encoder.build(input_batch_stream); + let encode_stream = encoder.build(input_batch_stream, ""); // use lower level stream to get access to app metadata - let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream).into_inner(); + let decode_stream = + FlightRecordBatchStream::new_from_flight_data(encode_stream, "").into_inner(); let messages: Vec<_> = decode_stream.try_collect().await.expect("encode fails"); @@ -275,9 +277,9 @@ async fn test_max_message_size_fuzz() { let input_batch_stream = futures::stream::iter(input.clone()).map(Ok); - let encode_stream = encoder.build(input_batch_stream); + let encode_stream = encoder.build(input_batch_stream, ""); - let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream); + let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream, ""); let output: Vec<_> = decode_stream.try_collect().await.expect("encode / decode"); for b in &output { @@ -299,7 +301,7 @@ async fn test_mismatched_record_batch_schema() { ]); let encoder = FlightDataEncoderBuilder::default(); - let encode_stream = encoder.build(input_batch_stream); + let encode_stream = encoder.build(input_batch_stream, ""); let result: Result, FlightError> = encode_stream.try_collect().await; let err = result.unwrap_err(); @@ -315,16 +317,16 @@ async fn test_chained_streams_batch_decoder() { let batch2 = make_dictionary_batch(3); // Model sending two flight streams back to back, with different schemas - let encode_stream1 = - FlightDataEncoderBuilder::default().build(futures::stream::iter(vec![Ok(batch1.clone())])); - let encode_stream2 = - FlightDataEncoderBuilder::default().build(futures::stream::iter(vec![Ok(batch2.clone())])); + let encode_stream1 = FlightDataEncoderBuilder::default() + .build(futures::stream::iter(vec![Ok(batch1.clone())]), ""); + let encode_stream2 = FlightDataEncoderBuilder::default() + .build(futures::stream::iter(vec![Ok(batch2.clone())]), ""); // append the two streams (so they will have two different schema messages) let encode_stream = encode_stream1.chain(encode_stream2); // FlightRecordBatchStream errors if the schema changes - let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream); + let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream, ""); let result: Result, FlightError> = decode_stream.try_collect().await; let err = result.unwrap_err(); @@ -340,16 +342,16 @@ async fn test_chained_streams_data_decoder() { let batch2 = make_dictionary_batch(3); // Model sending two flight streams back to back, with different schemas - let encode_stream1 = - FlightDataEncoderBuilder::default().build(futures::stream::iter(vec![Ok(batch1.clone())])); - let encode_stream2 = - FlightDataEncoderBuilder::default().build(futures::stream::iter(vec![Ok(batch2.clone())])); + let encode_stream1 = FlightDataEncoderBuilder::default() + .build(futures::stream::iter(vec![Ok(batch1.clone())]), ""); + let encode_stream2 = FlightDataEncoderBuilder::default() + .build(futures::stream::iter(vec![Ok(batch2.clone())]), ""); // append the two streams (so they will have two different schema messages) let encode_stream = encode_stream1.chain(encode_stream2); // lower level decode stream can handle multiple schema messages - let decode_stream = FlightDataDecoder::new(encode_stream); + let decode_stream = FlightDataDecoder::new(encode_stream, ""); let decoded_data: Vec<_> = decode_stream.try_collect().await.expect("encode / decode"); @@ -375,11 +377,11 @@ async fn test_mismatched_schema_message() { // and expect an error async fn do_test(batch1: RecordBatch, batch2: RecordBatch, expected: &str) { let encode_stream1 = FlightDataEncoderBuilder::default() - .build(futures::stream::iter(vec![Ok(batch1.clone())])) + .build(futures::stream::iter(vec![Ok(batch1.clone())]), "") // take only schema message from first stream .take(1); let encode_stream2 = FlightDataEncoderBuilder::default() - .build(futures::stream::iter(vec![Ok(batch2.clone())])) + .build(futures::stream::iter(vec![Ok(batch2.clone())]), "") // take only data message from second .skip(1); @@ -387,7 +389,7 @@ async fn test_mismatched_schema_message() { let encode_stream = encode_stream1.chain(encode_stream2); // FlightRecordBatchStream errors if the schema changes - let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream); + let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream, ""); let result: Result, FlightError> = decode_stream.try_collect().await; let err = result.unwrap_err().to_string(); @@ -446,9 +448,9 @@ async fn roundtrip_with_encoder( let input_batch_stream = futures::stream::iter(input_batches.clone()).map(Ok); - let encode_stream = encoder.build(input_batch_stream); + let encode_stream = encoder.build(input_batch_stream, ""); - let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream); + let decode_stream = FlightRecordBatchStream::new_from_flight_data(encode_stream, ""); let output_batches: Vec<_> = decode_stream.try_collect().await.expect("encode / decode"); // remove any empty batches from input as they are not transmitted diff --git a/arrow-flight/tests/flight_sql_client.rs b/arrow-flight/tests/flight_sql_client.rs index 349da062a82d..debbde43ccc8 100644 --- a/arrow-flight/tests/flight_sql_client.rs +++ b/arrow-flight/tests/flight_sql_client.rs @@ -206,6 +206,7 @@ impl FlightSqlService for FlightSqlServiceImpl { ) -> Result { let batches: Vec = FlightRecordBatchStream::new_from_flight_data( request.into_inner().map_err(|e| e.into()), + "", ) .try_collect() .await?; diff --git a/arrow-flight/tests/flight_sql_client_cli.rs b/arrow-flight/tests/flight_sql_client_cli.rs index c8e9190e246f..ef79c87a317b 100644 --- a/arrow-flight/tests/flight_sql_client_cli.rs +++ b/arrow-flight/tests/flight_sql_client_cli.rs @@ -619,7 +619,7 @@ impl FlightSqlService for FlightSqlServiceImpl { let batch = builder.build(); let stream = FlightDataEncoderBuilder::new() .with_schema(schema) - .build(futures::stream::once(async { batch })) + .build(futures::stream::once(async { batch }), "") .map_err(Status::from); Ok(Response::new(Box::pin(stream))) } @@ -642,7 +642,7 @@ impl FlightSqlService for FlightSqlServiceImpl { let batch = builder.build(); let stream = FlightDataEncoderBuilder::new() .with_schema(schema) - .build(futures::stream::once(async { batch })) + .build(futures::stream::once(async { batch }), "") .map_err(Status::from); Ok(Response::new(Box::pin(stream))) } @@ -685,7 +685,7 @@ impl FlightSqlService for FlightSqlServiceImpl { let batch = builder.build(); let stream = FlightDataEncoderBuilder::new() .with_schema(schema) - .build(futures::stream::once(async { batch })) + .build(futures::stream::once(async { batch }), "") .map_err(Status::from); Ok(Response::new(Box::pin(stream))) } @@ -704,7 +704,7 @@ impl FlightSqlService for FlightSqlServiceImpl { let batch = builder.build(); let stream = FlightDataEncoderBuilder::new() .with_schema(schema) - .build(futures::stream::once(async { batch })) + .build(futures::stream::once(async { batch }), "") .map_err(Status::from); Ok(Response::new(Box::pin(stream))) } @@ -717,6 +717,7 @@ impl FlightSqlService for FlightSqlServiceImpl { // just make sure decoding the parameters works let parameters = FlightRecordBatchStream::new_from_flight_data( request.into_inner().map_err(|e| e.into()), + "", ) .try_collect::>() .await?; From 60ef28553d3cb0f1d575a007456d8592d30aa669 Mon Sep 17 00:00:00 2001 From: Faiaz Sanaulla Date: Thu, 8 May 2025 13:19:15 +0200 Subject: [PATCH 2/5] update --- arrow-flight/Cargo.toml | 1 - arrow-flight/src/decode.rs | 9 ++++----- arrow-flight/src/encode.rs | 20 +++++--------------- 3 files changed, 9 insertions(+), 21 deletions(-) diff --git a/arrow-flight/Cargo.toml b/arrow-flight/Cargo.toml index 45880d644c6f..9531c63af003 100644 --- a/arrow-flight/Cargo.toml +++ b/arrow-flight/Cargo.toml @@ -49,7 +49,6 @@ prost = { version = "0.13.1", default-features = false, features = ["prost-deriv prost-types = { version = "0.13.1", default-features = false } tokio = { version = "1.0", default-features = false, features = ["macros", "rt", "rt-multi-thread"] } tonic = { version = "0.13", default-features = false, features = ["transport", "codegen", "prost"] } -tracing = { version = "0.1.40", features = ["default", "log"] } # CLI-related dependencies anyhow = { version = "1.0", optional = true } diff --git a/arrow-flight/src/decode.rs b/arrow-flight/src/decode.rs index 7dae90d000ca..7cc101fde19f 100644 --- a/arrow-flight/src/decode.rs +++ b/arrow-flight/src/decode.rs @@ -23,7 +23,6 @@ use bytes::Bytes; use futures::{ready, stream::BoxStream, Stream, StreamExt}; use std::{collections::HashMap, fmt::Debug, pin::Pin, sync::Arc, task::Poll}; use tonic::metadata::MetadataMap; -use tracing::debug; use crate::error::{FlightError, Result}; @@ -358,24 +357,24 @@ impl futures::Stream for FlightDataDecoder { cx: &mut std::task::Context<'_>, ) -> Poll> { if self.done { - debug!(self.reader_id, "stream done"); + println!("stream done, {}", self.reader_id); return Poll::Ready(None); } loop { - debug!(self.reader_id, "polling next message"); + println!("polling next message, {}", self.reader_id); let res = ready!(self.response.poll_next_unpin(cx)); return Poll::Ready(match res { None => { self.done = true; - debug!(self.reader_id, "inner is exhausted"); + println!("inner is exhausted, {}", self.reader_id); None // inner is exhausted } Some(data) => Some(match data { Err(e) => Err(e), Ok(data) => match self.extract_message(data) { Ok(Some(extracted)) => { - debug!(self.reader_id, "message extracted"); + println!("message extracted, {}", self.reader_id); Ok(extracted) } Ok(None) => continue, // Need next input message diff --git a/arrow-flight/src/encode.rs b/arrow-flight/src/encode.rs index b7ba68e1f870..9a9165e55116 100644 --- a/arrow-flight/src/encode.rs +++ b/arrow-flight/src/encode.rs @@ -25,7 +25,6 @@ use arrow_ipc::writer::{DictionaryTracker, IpcDataGenerator, IpcWriteOptions}; use arrow_schema::{DataType, Field, FieldRef, Fields, Schema, SchemaRef, UnionMode}; use bytes::Bytes; use futures::{ready, stream::BoxStream, Stream, StreamExt}; -use tracing::debug; /// Creates a [`Stream`] of [`FlightData`]s from a /// `Stream` of [`Result`]<[`RecordBatch`], [`FlightError`]>. @@ -405,21 +404,15 @@ impl Stream for FlightDataEncoder { cx: &mut std::task::Context<'_>, ) -> Poll> { loop { - debug!( - reader_id = self.reader_id, - "flight data encoder polling next" - ); + println!("flight data encoder polling next, {}", self.reader_id); if self.done && self.queue.is_empty() { - debug!( - reader_id = self.reader_id, - "stream done, no more data to send" - ); + println!("stream done, no more data to send, {}", self.reader_id); return Poll::Ready(None); } // Any messages queued to send? if let Some(data) = self.queue.pop_front() { - debug!(reader_id = self.reader_id, "sending queued message"); + println!("sending queued message, {}", self.reader_id); return Poll::Ready(Some(Ok(data))); } @@ -432,10 +425,7 @@ impl Stream for FlightDataEncoder { self.done = true; // queue must also be empty so we are done assert!(self.queue.is_empty()); - debug!( - reader_id = self.reader_id, - "stream done, no more data to send" - ); + println!("stream done, no more data to send, {}", self.reader_id); return Poll::Ready(None); } Some(Err(e)) => { @@ -445,7 +435,7 @@ impl Stream for FlightDataEncoder { return Poll::Ready(Some(Err(e))); } Some(Ok(batch)) => { - debug!(reader_id = self.reader_id, "got batch"); + println!("got batch, {}", self.reader_id); // had data, encode into the queue if let Err(e) = self.encode_batch(batch) { self.done = true; From be6eac6a3ce65fbf7a3e2267b5e5640101d0db54 Mon Sep 17 00:00:00 2001 From: Faiaz Sanaulla Date: Thu, 8 May 2025 14:25:36 +0200 Subject: [PATCH 3/5] poll count --- arrow-flight/src/decode.rs | 24 ++++++++++++++++++++---- arrow-flight/src/encode.rs | 25 +++++++++++++++++++++---- 2 files changed, 41 insertions(+), 8 deletions(-) diff --git a/arrow-flight/src/decode.rs b/arrow-flight/src/decode.rs index 7cc101fde19f..dea9c36c9906 100644 --- a/arrow-flight/src/decode.rs +++ b/arrow-flight/src/decode.rs @@ -236,6 +236,8 @@ pub struct FlightDataDecoder { done: bool, reader_id: String, + + poll_count: usize, } impl Debug for FlightDataDecoder { @@ -259,6 +261,7 @@ impl FlightDataDecoder { response: response.boxed(), done: false, reader_id: reader_id.to_string(), + poll_count: 0, } } @@ -357,24 +360,37 @@ impl futures::Stream for FlightDataDecoder { cx: &mut std::task::Context<'_>, ) -> Poll> { if self.done { - println!("stream done, {}", self.reader_id); + println!( + "flight data decoder - stream done, {}/{}", + self.reader_id, self.poll_count + ); return Poll::Ready(None); } + self.poll_count += 1; loop { - println!("polling next message, {}", self.reader_id); + println!( + "flight data decoder polling next, {}/{}", + self.reader_id, self.poll_count + ); let res = ready!(self.response.poll_next_unpin(cx)); return Poll::Ready(match res { None => { self.done = true; - println!("inner is exhausted, {}", self.reader_id); + println!( + "flight data decoder inner is exhausted, {}/{}", + self.reader_id, self.poll_count + ); None // inner is exhausted } Some(data) => Some(match data { Err(e) => Err(e), Ok(data) => match self.extract_message(data) { Ok(Some(extracted)) => { - println!("message extracted, {}", self.reader_id); + println!( + "flight data decoder message extracted, {}/{}", + self.reader_id, self.poll_count + ); Ok(extracted) } Ok(None) => continue, // Need next input message diff --git a/arrow-flight/src/encode.rs b/arrow-flight/src/encode.rs index 9a9165e55116..44e3036b44b3 100644 --- a/arrow-flight/src/encode.rs +++ b/arrow-flight/src/encode.rs @@ -290,6 +290,8 @@ pub struct FlightDataEncoder { dictionary_handling: DictionaryHandling, reader_id: String, + + poll_count: usize, } impl FlightDataEncoder { @@ -318,6 +320,7 @@ impl FlightDataEncoder { descriptor, dictionary_handling, reader_id: reader_id.to_string(), + poll_count: 0, }; // If schema is known up front, enqueue it immediately @@ -404,15 +407,23 @@ impl Stream for FlightDataEncoder { cx: &mut std::task::Context<'_>, ) -> Poll> { loop { + self.poll_count += 1; + println!("flight data encoder polling next, {}", self.reader_id); if self.done && self.queue.is_empty() { - println!("stream done, no more data to send, {}", self.reader_id); + println!( + "flight data encoder stream done, no more data to send, {}/{}", + self.reader_id, self.poll_count + ); return Poll::Ready(None); } // Any messages queued to send? if let Some(data) = self.queue.pop_front() { - println!("sending queued message, {}", self.reader_id); + println!( + "flight data encoder sending queued message, {}/{}", + self.reader_id, self.poll_count + ); return Poll::Ready(Some(Ok(data))); } @@ -425,7 +436,10 @@ impl Stream for FlightDataEncoder { self.done = true; // queue must also be empty so we are done assert!(self.queue.is_empty()); - println!("stream done, no more data to send, {}", self.reader_id); + println!( + "flight data encoder stream done, no more data to send, {}/{}", + self.reader_id, self.poll_count + ); return Poll::Ready(None); } Some(Err(e)) => { @@ -435,7 +449,10 @@ impl Stream for FlightDataEncoder { return Poll::Ready(Some(Err(e))); } Some(Ok(batch)) => { - println!("got batch, {}", self.reader_id); + println!( + "flight data encoder got batch, {}/{}", + self.reader_id, self.poll_count + ); // had data, encode into the queue if let Err(e) = self.encode_batch(batch) { self.done = true; From 70925c99187d6acdaa3472f598a66897d711117c Mon Sep 17 00:00:00 2001 From: Faiaz Sanaulla Date: Thu, 8 May 2025 15:24:31 +0200 Subject: [PATCH 4/5] missing poll count --- arrow-flight/src/encode.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/arrow-flight/src/encode.rs b/arrow-flight/src/encode.rs index 44e3036b44b3..ab6006c5de60 100644 --- a/arrow-flight/src/encode.rs +++ b/arrow-flight/src/encode.rs @@ -409,7 +409,7 @@ impl Stream for FlightDataEncoder { loop { self.poll_count += 1; - println!("flight data encoder polling next, {}", self.reader_id); + println!("flight data encoder polling next, {}/{}", self.reader_id, self.poll_count); if self.done && self.queue.is_empty() { println!( "flight data encoder stream done, no more data to send, {}/{}", From e979b1c10b9d90b8bfccf0630bab281ca1b64d69 Mon Sep 17 00:00:00 2001 From: Faiaz Sanaulla Date: Thu, 8 May 2025 15:41:27 +0200 Subject: [PATCH 5/5] measure poll time --- arrow-flight/src/decode.rs | 9 ++++++++- arrow-flight/src/encode.rs | 12 +++++++++++- 2 files changed, 19 insertions(+), 2 deletions(-) diff --git a/arrow-flight/src/decode.rs b/arrow-flight/src/decode.rs index dea9c36c9906..5adffd97e369 100644 --- a/arrow-flight/src/decode.rs +++ b/arrow-flight/src/decode.rs @@ -21,7 +21,7 @@ use arrow_buffer::Buffer; use arrow_schema::{Schema, SchemaRef}; use bytes::Bytes; use futures::{ready, stream::BoxStream, Stream, StreamExt}; -use std::{collections::HashMap, fmt::Debug, pin::Pin, sync::Arc, task::Poll}; +use std::{collections::HashMap, fmt::Debug, pin::Pin, sync::Arc, task::Poll, time::Instant}; use tonic::metadata::MetadataMap; use crate::error::{FlightError, Result}; @@ -372,7 +372,14 @@ impl futures::Stream for FlightDataDecoder { "flight data decoder polling next, {}/{}", self.reader_id, self.poll_count ); + let now = Instant::now(); let res = ready!(self.response.poll_next_unpin(cx)); + println!( + "flight data decoder polled next, {}/{} - took {:?}", + self.reader_id, + self.poll_count, + now.elapsed() + ); return Poll::Ready(match res { None => { diff --git a/arrow-flight/src/encode.rs b/arrow-flight/src/encode.rs index ab6006c5de60..f4397ab8b800 100644 --- a/arrow-flight/src/encode.rs +++ b/arrow-flight/src/encode.rs @@ -409,7 +409,10 @@ impl Stream for FlightDataEncoder { loop { self.poll_count += 1; - println!("flight data encoder polling next, {}/{}", self.reader_id, self.poll_count); + println!( + "flight data encoder polling next, {}/{}", + self.reader_id, self.poll_count + ); if self.done && self.queue.is_empty() { println!( "flight data encoder stream done, no more data to send, {}/{}", @@ -428,7 +431,14 @@ impl Stream for FlightDataEncoder { } // Get next batch + let now = std::time::Instant::now(); let batch = ready!(self.inner.poll_next_unpin(cx)); + println!( + "flight data encoder polled next, {}/{} - took {:?}", + self.reader_id, + self.poll_count, + now.elapsed() + ); match batch { None => {