From 959d82a1604b9252ad07764f747f1933f2c30dbf Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Fri, 3 Jul 2026 15:53:31 +0800 Subject: [PATCH 1/8] feat(arrow-ipc): add sans-IO stream encoder Introduce StreamEncoder for IPC streaming without requiring a std::io::Write sink. The encoder owns stream state, emits ordered Buffer chunks, and preserves the low-copy path for uncompressed record batch body buffers. Add byte-for-byte compatibility tests against StreamWriter for normal batches, empty streams, and dictionary batches. Closes #7812 Signed-off-by: Jiawei Zhao --- arrow-ipc/src/writer.rs | 334 ++++++++++++++++++++++++++++++++++++++++ 1 file changed, 334 insertions(+) diff --git a/arrow-ipc/src/writer.rs b/arrow-ipc/src/writer.rs index 311cb44e4665..c6ba0f8a015a 100644 --- a/arrow-ipc/src/writer.rs +++ b/arrow-ipc/src/writer.rs @@ -725,6 +725,67 @@ impl IpcDataGenerator { }) } + /// Encode dictionary batches and the record batch to output buffers, skipping the + /// intermediate body `Vec` allocations for uncompressed record batch buffers. + fn encode_to_buffers( + &self, + batch: &RecordBatch, + dictionary_tracker: &mut DictionaryTracker, + write_options: &IpcWriteOptions, + ipc_write_context: &mut IpcWriteContext, + out: &mut Vec, + ) -> Result { + let encoded_dictionaries = + self.encode_all_dicts(batch, dictionary_tracker, write_options, ipc_write_context)?; + + let mut dictionary_block_sizes = Vec::with_capacity(encoded_dictionaries.len()); + for dict in encoded_dictionaries { + dictionary_block_sizes.push(encoded_data_to_buffers(out, dict, write_options)?); + } + + let capacity = batch + .columns() + .iter() + .map(|a| estimate_encoded_buffer_count(a.data_type())) + .sum(); + let mut encoded_buffers: Vec = Vec::with_capacity(capacity); + let (ipc_message, body_len, tail_pad) = self.record_batch_to_bytes( + batch, + write_options, + ipc_write_context, + &mut IpcBodySink::Collect(&mut encoded_buffers), + )?; + + let alignment = write_options.alignment; + let alignment_mask = usize::from(alignment - 1); + let prefix_size = if write_options.write_legacy_ipc_format { + 4 + } else { + 8 + }; + let ipc_message_len = ipc_message.len(); + let aligned_size = (ipc_message_len + prefix_size + alignment_mask) & !alignment_mask; + push_continuation_buffer(out, write_options, (aligned_size - prefix_size) as i32)?; + out.push(Buffer::from(ipc_message)); + push_padding_buffer(out, aligned_size - ipc_message_len - prefix_size); + for enc in encoded_buffers { + let len = enc.len(); + let buffer = match enc { + EncodedBuffer::Raw(buffer) => buffer, + EncodedBuffer::Compressed(bytes) => Buffer::from(bytes), + }; + out.push(buffer); + push_padding_buffer(out, pad_to_alignment(alignment, len)); + } + push_padding_buffer(out, tail_pad); + + Ok(IpcWriteMetadata { + dictionary_block_sizes, + padded_header_len: aligned_size, + body_len, + }) + } + /// Encodes a batch to a number of [EncodedData] items (dictionary batches + the record batch). /// The [DictionaryTracker] keeps track of dictionaries with new `dict_id`s (so they are only sent once) /// Make sure the [DictionaryTracker] is initialized at the start of the stream. @@ -1564,6 +1625,139 @@ impl RecordBatchWriter for FileWriter { } } +/// Arrow IPC stream encoder. +/// +/// Encodes Arrow [`RecordBatch`]es to byte buffers using the [IPC Streaming Format], +/// without performing any IO. +/// +/// The returned [`Buffer`]s are ordered and should be written to the destination +/// stream in order. Uncompressed record batch body buffers can share the original +/// Arrow buffers instead of being copied into an intermediate contiguous buffer. +/// +/// # Example +/// ``` +/// # use arrow_array::record_batch; +/// # use arrow_ipc::writer::StreamEncoder; +/// # use arrow_schema::ArrowError; +/// # fn main() -> Result<(), ArrowError> { +/// let batch = record_batch!(("a", Int32, [1, 2, 3]))?; +/// +/// let mut encoder = StreamEncoder::try_new(&batch.schema())?; +/// let mut stream = vec![]; +/// for buffer in encoder.encode(&batch)? { +/// stream.extend_from_slice(buffer.as_slice()); +/// } +/// for buffer in encoder.finish()? { +/// stream.extend_from_slice(buffer.as_slice()); +/// } +/// # Ok(()) +/// # } +/// ``` +/// +/// [IPC Streaming Format]: https://arrow.apache.org/docs/format/Columnar.html#ipc-streaming-format +pub struct StreamEncoder { + schema: Schema, + /// IPC write options + write_options: IpcWriteOptions, + /// Whether the stream schema has been encoded + schema_encoded: bool, + /// Whether the end-of-stream marker has been encoded + finished: bool, + /// Keeps track of dictionaries that have been encoded + dictionary_tracker: DictionaryTracker, + data_gen: IpcDataGenerator, + ipc_write_context: IpcWriteContext, +} + +impl StreamEncoder { + /// Try to create a new stream encoder. + pub fn try_new(schema: &Schema) -> Result { + let write_options = IpcWriteOptions::default(); + Self::try_new_with_options(schema, write_options) + } + + /// Try to create a new stream encoder with [`IpcWriteOptions`]. + pub fn try_new_with_options( + schema: &Schema, + write_options: IpcWriteOptions, + ) -> Result { + ensure_supported_ipc_schema(schema)?; + + Ok(Self { + schema: schema.clone(), + write_options, + schema_encoded: false, + finished: false, + dictionary_tracker: DictionaryTracker::new(false), + data_gen: IpcDataGenerator::default(), + ipc_write_context: IpcWriteContext::default(), + }) + } + + /// Encode a [`RecordBatch`] into buffers. + /// + /// The first call also includes the IPC stream schema message before the + /// record batch message. Later calls only include dictionary and record + /// batch messages. + /// + /// # Errors + /// + /// Returns an error if the encoder is already finished or encoding fails. + pub fn encode(&mut self, batch: &RecordBatch) -> Result, ArrowError> { + if self.finished { + return Err(ArrowError::IpcError( + "Cannot encode record batch to stream encoder as it is closed".to_string(), + )); + } + + let mut out = vec![]; + self.encode_schema(&mut out)?; + self.data_gen.encode_to_buffers( + batch, + &mut self.dictionary_tracker, + &self.write_options, + &mut self.ipc_write_context, + &mut out, + )?; + Ok(out) + } + + /// Encode the end-of-stream marker and mark this encoder as finished. + /// + /// If no batches have been encoded, this also emits the IPC stream schema + /// message so the returned buffers form a valid empty IPC stream. + /// + /// # Errors + /// + /// Returns an error if the encoder is already finished. + pub fn finish(&mut self) -> Result, ArrowError> { + if self.finished { + return Err(ArrowError::IpcError( + "Cannot finish stream encoder as it is closed".to_string(), + )); + } + + let mut out = vec![]; + self.encode_schema(&mut out)?; + push_continuation_buffer(&mut out, &self.write_options, 0)?; + self.finished = true; + Ok(out) + } + + fn encode_schema(&mut self, out: &mut Vec) -> Result<(), ArrowError> { + if !self.schema_encoded { + let encoded_message = self.data_gen.schema_to_bytes_with_dictionary_tracker( + &self.schema, + &mut self.dictionary_tracker, + &self.write_options, + ); + encoded_data_to_buffers(out, encoded_message, &self.write_options)?; + self.schema_encoded = true; + } + Ok(()) + } +} + /// Arrow Stream Writer /// /// Writes Arrow [`RecordBatch`]es to bytes using the [IPC Streaming Format]. @@ -1823,6 +2017,63 @@ pub struct EncodedData { /// Arrow buffers to be written, should be an empty vec for schema messages pub arrow_data: Vec, } + +fn encoded_data_to_buffers( + out: &mut Vec, + encoded: EncodedData, + write_options: &IpcWriteOptions, +) -> Result<(usize, usize), ArrowError> { + let arrow_data_len = encoded.arrow_data.len(); + if arrow_data_len % usize::from(write_options.alignment) != 0 { + return Err(ArrowError::MemoryError( + "Arrow data not aligned".to_string(), + )); + } + + let alignment_mask = usize::from(write_options.alignment - 1); + let flatbuf_size = encoded.ipc_message.len(); + let prefix_size = if write_options.write_legacy_ipc_format { + 4 + } else { + 8 + }; + let aligned_size = (flatbuf_size + prefix_size + alignment_mask) & !alignment_mask; + let padding_bytes = aligned_size - flatbuf_size - prefix_size; + + push_continuation_buffer(out, write_options, (aligned_size - prefix_size) as i32)?; + if flatbuf_size > 0 { + out.push(Buffer::from(encoded.ipc_message)); + } + push_padding_buffer(out, padding_bytes); + + let body_len = if arrow_data_len > 0 { + out.push(Buffer::from(encoded.arrow_data)); + arrow_data_len + } else { + 0 + }; + + Ok((aligned_size, body_len)) +} + +fn push_continuation_buffer( + out: &mut Vec, + write_options: &IpcWriteOptions, + total_len: i32, +) -> Result<(), ArrowError> { + // Continuation bytes are generated stream framing, so they need a small owned buffer. + let mut buffer = Vec::with_capacity(8); + write_continuation(&mut buffer, write_options, total_len)?; + out.push(Buffer::from(buffer)); + Ok(()) +} + +fn push_padding_buffer(out: &mut Vec, len: usize) { + if len > 0 { + out.push(Buffer::from(&PADDING[..len])); + } +} + /// Write a message's IPC data and buffers, returning metadata and buffer data lengths written pub fn write_message( mut writer: W, @@ -2523,6 +2774,89 @@ mod tests { stream_reader.next().unwrap().unwrap() } + fn encode_stream( + schema: &Schema, + batches: &[RecordBatch], + options: IpcWriteOptions, + ) -> Vec { + let mut encoder = StreamEncoder::try_new_with_options(schema, options).unwrap(); + let mut bytes = Vec::new(); + for batch in batches { + for buffer in encoder.encode(batch).unwrap() { + bytes.extend_from_slice(buffer.as_slice()); + } + } + for buffer in encoder.finish().unwrap() { + bytes.extend_from_slice(buffer.as_slice()); + } + bytes + } + + fn write_stream(schema: &Schema, batches: &[RecordBatch], options: IpcWriteOptions) -> Vec { + let mut bytes = Vec::new(); + let mut writer = StreamWriter::try_new_with_options(&mut bytes, schema, options).unwrap(); + for batch in batches { + writer.write(batch).unwrap(); + } + writer.finish().unwrap(); + bytes + } + + // StreamEncoder and StreamWriter currently use separate encoding paths, so these tests + // verify the new sans-IO API preserves the existing IPC stream byte layout. + #[test] + fn test_stream_encoder_matches_stream_writer() { + let batch = record_batch!(("a", Int32, [1, 2, 3]), ("b", Utf8, ["x", "y", "z"])).unwrap(); + let options = IpcWriteOptions::default(); + let encoded = encode_stream(batch.schema_ref(), &[batch.clone()], options.clone()); + let written = write_stream(batch.schema_ref(), &[batch.clone()], options); + + assert_eq!(encoded, written); + + let mut reader = StreamReader::try_new(Cursor::new(encoded), None).unwrap(); + assert_eq!(reader.next().unwrap().unwrap(), batch); + assert!(reader.next().is_none()); + } + + #[test] + fn test_stream_encoder_empty_stream_matches_stream_writer() { + let schema = Schema::new(vec![Field::new("a", DataType::Int32, true)]); + let options = IpcWriteOptions::default(); + let encoded = encode_stream(&schema, &[], options.clone()); + let written = write_stream(&schema, &[], options); + + assert_eq!(encoded, written); + + let mut reader = StreamReader::try_new(Cursor::new(encoded), None).unwrap(); + assert!(reader.next().is_none()); + } + + #[test] + fn test_stream_encoder_dictionary_batches_match_stream_writer() { + let schema = Arc::new(Schema::new(vec![Field::new( + "a", + DataType::Dictionary(Box::new(DataType::UInt8), Box::new(DataType::Utf8)), + false, + )])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![Arc::new(DictionaryArray::new( + UInt8Array::from_iter_values([0, 1, 0]), + Arc::new(StringArray::from_iter_values(["a", "b"])), + ))], + ) + .unwrap(); + let options = IpcWriteOptions::default(); + let encoded = encode_stream(&schema, &[batch.clone()], options.clone()); + let written = write_stream(&schema, &[batch.clone()], options); + + assert_eq!(encoded, written); + + let mut reader = StreamReader::try_new(Cursor::new(encoded), None).unwrap(); + assert_eq!(reader.next().unwrap().unwrap(), batch); + assert!(reader.next().is_none()); + } + #[test] #[cfg(feature = "lz4")] fn test_write_empty_record_batch_lz4_compression() { From 0405bcf07bb0d334b0253f2de7c5c72bf4602a35 Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Mon, 13 Jul 2026 13:23:06 +0800 Subject: [PATCH 2/8] test(arrow-ipc): cover async stream encoding Exercise StreamEncoder with a Tokio AsyncWrite flow and compare the bytes with StreamWriter. This keeps the PR tied to the async writer use case from the original issue. Refs #7812 Signed-off-by: Jiawei Zhao --- arrow-ipc/Cargo.toml | 2 +- arrow-ipc/src/writer.rs | 38 +++++++++++++++++++++++++++++--------- 2 files changed, 30 insertions(+), 10 deletions(-) diff --git a/arrow-ipc/Cargo.toml b/arrow-ipc/Cargo.toml index 601eb86b288d..e8c821145b94 100644 --- a/arrow-ipc/Cargo.toml +++ b/arrow-ipc/Cargo.toml @@ -52,7 +52,7 @@ lz4 = ["lz4_flex"] [dev-dependencies] criterion = { workspace = true } tempfile = "3.3" -tokio = "1.43.0" +tokio = { version = "1.43.0", features = ["io-util", "macros", "rt"] } # used in benches memmap2 = "0.9.3" bytes = "1.9" diff --git a/arrow-ipc/src/writer.rs b/arrow-ipc/src/writer.rs index c6ba0f8a015a..860d5e3ea600 100644 --- a/arrow-ipc/src/writer.rs +++ b/arrow-ipc/src/writer.rs @@ -1653,8 +1653,6 @@ impl RecordBatchWriter for FileWriter { /// # Ok(()) /// # } /// ``` -/// -/// [IPC Streaming Format]: https://arrow.apache.org/docs/format/Columnar.html#ipc-streaming-format pub struct StreamEncoder { schema: Schema, /// IPC write options @@ -2804,14 +2802,36 @@ mod tests { // StreamEncoder and StreamWriter currently use separate encoding paths, so these tests // verify the new sans-IO API preserves the existing IPC stream byte layout. - #[test] - fn test_stream_encoder_matches_stream_writer() { + #[tokio::test] + async fn test_stream_encoder_async_writer_matches_stream_writer() { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + let batch = record_batch!(("a", Int32, [1, 2, 3]), ("b", Utf8, ["x", "y", "z"])).unwrap(); let options = IpcWriteOptions::default(); - let encoded = encode_stream(batch.schema_ref(), &[batch.clone()], options.clone()); - let written = write_stream(batch.schema_ref(), &[batch.clone()], options); + let expected = write_stream( + batch.schema_ref(), + std::slice::from_ref(&batch), + options.clone(), + ); - assert_eq!(encoded, written); + let (mut sink, mut source) = tokio::io::duplex(64); + let read = tokio::spawn(async move { + let mut bytes = Vec::new(); + source.read_to_end(&mut bytes).await.unwrap(); + bytes + }); + + let mut encoder = StreamEncoder::try_new_with_options(batch.schema_ref(), options).unwrap(); + for buffer in encoder.encode(&batch).unwrap() { + sink.write_all(buffer.as_slice()).await.unwrap(); + } + for buffer in encoder.finish().unwrap() { + sink.write_all(buffer.as_slice()).await.unwrap(); + } + sink.shutdown().await.unwrap(); + + let encoded = read.await.unwrap(); + assert_eq!(encoded, expected); let mut reader = StreamReader::try_new(Cursor::new(encoded), None).unwrap(); assert_eq!(reader.next().unwrap().unwrap(), batch); @@ -2847,8 +2867,8 @@ mod tests { ) .unwrap(); let options = IpcWriteOptions::default(); - let encoded = encode_stream(&schema, &[batch.clone()], options.clone()); - let written = write_stream(&schema, &[batch.clone()], options); + let encoded = encode_stream(&schema, std::slice::from_ref(&batch), options.clone()); + let written = write_stream(&schema, std::slice::from_ref(&batch), options); assert_eq!(encoded, written); From b80b9dc07587d5100aa6b0d07958fc7db8a30ccf Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Fri, 17 Jul 2026 09:17:33 +0800 Subject: [PATCH 3/8] refactor(ipc): share message sink Route writer and buffer output through a shared internal sink. This avoids duplicated IPC framing while preserving Buffer output. Signed-off-by: Jiawei Zhao --- arrow-ipc/src/writer.rs | 447 +++++++++++++++++++--------------------- 1 file changed, 217 insertions(+), 230 deletions(-) diff --git a/arrow-ipc/src/writer.rs b/arrow-ipc/src/writer.rs index 860d5e3ea600..4e2ad4e5c7a6 100644 --- a/arrow-ipc/src/writer.rs +++ b/arrow-ipc/src/writer.rs @@ -135,6 +135,165 @@ impl<'a> IpcBodySink<'a> { } } +/// Destination for a complete framed IPC message. +/// +/// This emits the stream/file framing around the serialized FlatBuffer +/// [`crate::Message`] metadata plus its optional body buffers. +enum IpcMessageSink<'a> { + /// Write bytes directly to a synchronous writer. + Writer(&'a mut dyn Write), + /// Accumulate ordered buffers for deferred writing. + Buffers(&'a mut Vec), +} + +impl IpcMessageSink<'_> { + fn write_vec(&mut self, bytes: Vec) -> Result<(), ArrowError> { + if bytes.is_empty() { + return Ok(()); + } + + match self { + Self::Writer(writer) => writer.write_all(&bytes)?, + Self::Buffers(out) => out.push(Buffer::from(bytes)), + } + Ok(()) + } + + fn write_encoded_buffer(&mut self, buffer: EncodedBuffer) -> Result<(), ArrowError> { + match (self, buffer) { + (Self::Writer(writer), buffer) => writer.write_all(buffer.as_slice())?, + (Self::Buffers(out), EncodedBuffer::Raw(buffer)) => out.push(buffer), + (Self::Buffers(out), EncodedBuffer::Compressed(bytes)) => out.push(Buffer::from(bytes)), + } + Ok(()) + } + + fn write_padding(&mut self, len: usize) -> Result<(), ArrowError> { + if len == 0 { + return Ok(()); + } + + match self { + Self::Writer(writer) => writer.write_all(&PADDING[..len])?, + Self::Buffers(out) => out.push(Buffer::from(&PADDING[..len])), + } + Ok(()) + } + + /// Writes the IPC continuation marker and metadata length prefix. + fn write_continuation( + &mut self, + write_options: &IpcWriteOptions, + metadata_len: i32, + ) -> Result<(), ArrowError> { + let mut buffer = Vec::with_capacity(8); + + match write_options.metadata_version { + crate::MetadataVersion::V1 + | crate::MetadataVersion::V2 + | crate::MetadataVersion::V3 => { + unreachable!("Options with the metadata version cannot be created") + } + crate::MetadataVersion::V4 => { + if !write_options.write_legacy_ipc_format { + // v0.15.0 format + buffer.extend_from_slice(&CONTINUATION_MARKER); + } + buffer.extend_from_slice(&metadata_len.to_le_bytes()[..]); + } + crate::MetadataVersion::V5 => { + buffer.extend_from_slice(&CONTINUATION_MARKER); + buffer.extend_from_slice(&metadata_len.to_le_bytes()[..]); + } + z => panic!("Unsupported crate::MetadataVersion {z:?}"), + } + + self.write_vec(buffer) + } + + fn write_body_data(&mut self, data: Vec, alignment: u8) -> Result { + let len = data.len(); + let pad_len = pad_to_alignment(alignment, len); + self.write_vec(data)?; + self.write_padding(pad_len)?; + Ok(len + pad_len) + } + + fn write_encoded_data( + &mut self, + encoded: EncodedData, + write_options: &IpcWriteOptions, + ) -> Result<(usize, usize), ArrowError> { + let arrow_data_len = encoded.arrow_data.len(); + if arrow_data_len % usize::from(write_options.alignment) != 0 { + return Err(ArrowError::MemoryError( + "Arrow data not aligned".to_string(), + )); + } + + let alignment_mask = usize::from(write_options.alignment - 1); + let metadata = encoded.ipc_message; + let metadata_len = metadata.len(); + let prefix_size = if write_options.write_legacy_ipc_format { + 4 + } else { + 8 + }; + let padded_header_len = (metadata_len + prefix_size + alignment_mask) & !alignment_mask; + let padded_metadata_len = padded_header_len - prefix_size; + let metadata_padding = padded_metadata_len - metadata_len; + + self.write_continuation(write_options, padded_metadata_len as i32)?; + self.write_vec(metadata)?; + self.write_padding(metadata_padding)?; + + let body_len = if arrow_data_len > 0 { + self.write_body_data(encoded.arrow_data, write_options.alignment)? + } else { + 0 + }; + + Ok((padded_header_len, body_len)) + } + + fn write_record_batch( + &mut self, + metadata: Vec, + encoded_buffers: Vec, + body_len: usize, + tail_pad: usize, + write_options: &IpcWriteOptions, + ) -> Result<(usize, usize), ArrowError> { + let alignment = write_options.alignment; + let prefix_size = if write_options.write_legacy_ipc_format { + 4 + } else { + 8 + }; + let alignment_mask = usize::from(alignment - 1); + let padded_header_len = (metadata.len() + prefix_size + alignment_mask) & !alignment_mask; + let padded_metadata_len = padded_header_len - prefix_size; + let metadata_padding = padded_metadata_len - metadata.len(); + + self.write_continuation(write_options, padded_metadata_len as i32)?; + self.write_vec(metadata)?; + self.write_padding(metadata_padding)?; + for enc in encoded_buffers { + let len = enc.len(); + self.write_encoded_buffer(enc)?; + self.write_padding(pad_to_alignment(alignment, len))?; + } + self.write_padding(tail_pad)?; + + Ok((padded_header_len, body_len)) + } + + fn write_eos(&mut self, write_options: &IpcWriteOptions) -> Result<(), ArrowError> { + self.write_continuation(write_options, 0)?; + Ok(()) + } +} + /// Per-message sizes produced by [`IpcDataGenerator::write`]. /// /// [`FileWriter`] uses these to build the Block index entries required by the IPC footer for @@ -356,12 +515,12 @@ impl IpcDataGenerator { message.add_bodyLength(0); message.add_header(schema); // TODO: custom metadata - let data = message.finish(); - fbb.finish(data, None); + let root = message.finish(); + fbb.finish(root, None); - let data = fbb.finished_data(); + let metadata = fbb.finished_data(); EncodedData { - ipc_message: data.to_vec(), + ipc_message: metadata.to_vec(), arrow_data: vec![], } } @@ -623,7 +782,7 @@ impl IpcDataGenerator { let encoded_dictionaries = self.encode_all_dicts(batch, dictionary_tracker, write_options, ipc_write_context)?; let mut arrow_data = ipc_write_context.scratch(); - let (ipc_message, _, tail_pad) = self.record_batch_to_bytes( + let (metadata, _, tail_pad) = self.record_batch_to_bytes( batch, write_options, ipc_write_context, @@ -634,7 +793,7 @@ impl IpcDataGenerator { Ok(( encoded_dictionaries, EncodedData { - ipc_message, + ipc_message: metadata, arrow_data, }, )) @@ -676,53 +835,14 @@ impl IpcDataGenerator { ipc_write_context: &mut IpcWriteContext, writer: &mut W, ) -> Result { - let encoded_dictionaries = - self.encode_all_dicts(batch, dictionary_tracker, write_options, ipc_write_context)?; - - let mut dictionary_block_sizes = Vec::with_capacity(encoded_dictionaries.len()); - for dict in encoded_dictionaries { - dictionary_block_sizes.push(write_message(&mut *writer, dict, write_options)?); - } - - let capacity = batch - .columns() - .iter() - .map(|a| estimate_encoded_buffer_count(a.data_type())) - .sum(); - let mut encoded_buffers: Vec = Vec::with_capacity(capacity); - let (ipc_message, body_len, tail_pad) = self.record_batch_to_bytes( + let mut sink = IpcMessageSink::Writer(writer); + self.write_to_sink( batch, + dictionary_tracker, write_options, ipc_write_context, - &mut IpcBodySink::Collect(&mut encoded_buffers), - )?; - - let alignment = write_options.alignment; - let a = usize::from(alignment - 1); - let prefix_size = if write_options.write_legacy_ipc_format { - 4 - } else { - 8 - }; - let aligned_size = (ipc_message.len() + prefix_size + a) & !a; - write_continuation( - &mut *writer, - write_options, - (aligned_size - prefix_size) as i32, - )?; - writer.write_all(&ipc_message)?; - writer.write_all(&PADDING[..aligned_size - ipc_message.len() - prefix_size])?; - for enc in &encoded_buffers { - writer.write_all(enc.as_slice())?; - writer.write_all(&PADDING[..pad_to_alignment(alignment, enc.len())])?; - } - writer.write_all(&PADDING[..tail_pad])?; - - Ok(IpcWriteMetadata { - dictionary_block_sizes, - padded_header_len: aligned_size, - body_len, - }) + &mut sink, + ) } /// Encode dictionary batches and the record batch to output buffers, skipping the @@ -734,13 +854,31 @@ impl IpcDataGenerator { write_options: &IpcWriteOptions, ipc_write_context: &mut IpcWriteContext, out: &mut Vec, + ) -> Result { + let mut sink = IpcMessageSink::Buffers(out); + self.write_to_sink( + batch, + dictionary_tracker, + write_options, + ipc_write_context, + &mut sink, + ) + } + + fn write_to_sink( + &self, + batch: &RecordBatch, + dictionary_tracker: &mut DictionaryTracker, + write_options: &IpcWriteOptions, + ipc_write_context: &mut IpcWriteContext, + sink: &mut IpcMessageSink<'_>, ) -> Result { let encoded_dictionaries = self.encode_all_dicts(batch, dictionary_tracker, write_options, ipc_write_context)?; let mut dictionary_block_sizes = Vec::with_capacity(encoded_dictionaries.len()); for dict in encoded_dictionaries { - dictionary_block_sizes.push(encoded_data_to_buffers(out, dict, write_options)?); + dictionary_block_sizes.push(sink.write_encoded_data(dict, write_options)?); } let capacity = batch @@ -749,39 +887,19 @@ impl IpcDataGenerator { .map(|a| estimate_encoded_buffer_count(a.data_type())) .sum(); let mut encoded_buffers: Vec = Vec::with_capacity(capacity); - let (ipc_message, body_len, tail_pad) = self.record_batch_to_bytes( + let (metadata, body_len, tail_pad) = self.record_batch_to_bytes( batch, write_options, ipc_write_context, &mut IpcBodySink::Collect(&mut encoded_buffers), )?; - let alignment = write_options.alignment; - let alignment_mask = usize::from(alignment - 1); - let prefix_size = if write_options.write_legacy_ipc_format { - 4 - } else { - 8 - }; - let ipc_message_len = ipc_message.len(); - let aligned_size = (ipc_message_len + prefix_size + alignment_mask) & !alignment_mask; - push_continuation_buffer(out, write_options, (aligned_size - prefix_size) as i32)?; - out.push(Buffer::from(ipc_message)); - push_padding_buffer(out, aligned_size - ipc_message_len - prefix_size); - for enc in encoded_buffers { - let len = enc.len(); - let buffer = match enc { - EncodedBuffer::Raw(buffer) => buffer, - EncodedBuffer::Compressed(bytes) => Buffer::from(bytes), - }; - out.push(buffer); - push_padding_buffer(out, pad_to_alignment(alignment, len)); - } - push_padding_buffer(out, tail_pad); + let (padded_header_len, body_len) = + sink.write_record_batch(metadata, encoded_buffers, body_len, tail_pad, write_options)?; Ok(IpcWriteMetadata { dictionary_block_sizes, - padded_header_len: aligned_size, + padded_header_len, body_len, }) } @@ -807,7 +925,7 @@ impl IpcDataGenerator { /// Encodes a `RecordBatch` into a flatbuffer IPC message and fills `sink` with the /// serialised buffer data. /// - /// Returns `(ipc_message, body_len, tail_pad)`: the flatbuffer header bytes, the + /// Returns `(metadata, body_len, tail_pad)`: the FlatBuffer [`crate::Message`] bytes, the /// total body length including trailing padding, and the trailing alignment padding byte count. fn record_batch_to_bytes( &self, @@ -888,9 +1006,9 @@ impl IpcDataGenerator { let root = message.finish(); fbb.finish(root, None); - let ipc_message = fbb.finished_data().to_vec(); + let metadata = fbb.finished_data().to_vec(); fbb.reset(); - Ok((ipc_message, body_len, tail_pad)) + Ok((metadata, body_len, tail_pad)) } /// Write dictionary values into two sets of bytes, one for the header (crate::Message) and the @@ -988,11 +1106,11 @@ impl IpcDataGenerator { }; fbb.finish(root, None); - let ipc_message = fbb.finished_data().to_vec(); + let metadata = fbb.finished_data().to_vec(); fbb.reset(); Ok(EncodedData { - ipc_message, + ipc_message: metadata, arrow_data, }) } @@ -1536,7 +1654,10 @@ impl FileWriter { } // write EOS - write_continuation(&mut self.writer, &self.write_options, 0)?; + { + let mut sink = IpcMessageSink::Writer(&mut self.writer); + sink.write_eos(&self.write_options)?; + } let mut fbb = FlatBufferBuilder::new(); let dictionaries = fbb.create_vector(&self.dictionary_blocks); @@ -1737,7 +1858,8 @@ impl StreamEncoder { let mut out = vec![]; self.encode_schema(&mut out)?; - push_continuation_buffer(&mut out, &self.write_options, 0)?; + let mut sink = IpcMessageSink::Buffers(&mut out); + sink.write_eos(&self.write_options)?; self.finished = true; Ok(out) } @@ -1749,7 +1871,8 @@ impl StreamEncoder { &mut self.dictionary_tracker, &self.write_options, ); - encoded_data_to_buffers(out, encoded_message, &self.write_options)?; + let mut sink = IpcMessageSink::Buffers(out); + sink.write_encoded_data(encoded_message, &self.write_options)?; self.schema_encoded = true; } Ok(()) @@ -1924,7 +2047,10 @@ impl StreamWriter { )); } - write_continuation(&mut self.writer, &self.write_options, 0)?; + { + let mut sink = IpcMessageSink::Writer(&mut self.writer); + sink.write_eos(&self.write_options)?; + } self.writer.flush()?; self.finished = true; @@ -2016,158 +2142,14 @@ pub struct EncodedData { pub arrow_data: Vec, } -fn encoded_data_to_buffers( - out: &mut Vec, - encoded: EncodedData, - write_options: &IpcWriteOptions, -) -> Result<(usize, usize), ArrowError> { - let arrow_data_len = encoded.arrow_data.len(); - if arrow_data_len % usize::from(write_options.alignment) != 0 { - return Err(ArrowError::MemoryError( - "Arrow data not aligned".to_string(), - )); - } - - let alignment_mask = usize::from(write_options.alignment - 1); - let flatbuf_size = encoded.ipc_message.len(); - let prefix_size = if write_options.write_legacy_ipc_format { - 4 - } else { - 8 - }; - let aligned_size = (flatbuf_size + prefix_size + alignment_mask) & !alignment_mask; - let padding_bytes = aligned_size - flatbuf_size - prefix_size; - - push_continuation_buffer(out, write_options, (aligned_size - prefix_size) as i32)?; - if flatbuf_size > 0 { - out.push(Buffer::from(encoded.ipc_message)); - } - push_padding_buffer(out, padding_bytes); - - let body_len = if arrow_data_len > 0 { - out.push(Buffer::from(encoded.arrow_data)); - arrow_data_len - } else { - 0 - }; - - Ok((aligned_size, body_len)) -} - -fn push_continuation_buffer( - out: &mut Vec, - write_options: &IpcWriteOptions, - total_len: i32, -) -> Result<(), ArrowError> { - // Continuation bytes are generated stream framing, so they need a small owned buffer. - let mut buffer = Vec::with_capacity(8); - write_continuation(&mut buffer, write_options, total_len)?; - out.push(Buffer::from(buffer)); - Ok(()) -} - -fn push_padding_buffer(out: &mut Vec, len: usize) { - if len > 0 { - out.push(Buffer::from(&PADDING[..len])); - } -} - /// Write a message's IPC data and buffers, returning metadata and buffer data lengths written pub fn write_message( mut writer: W, encoded: EncodedData, write_options: &IpcWriteOptions, ) -> Result<(usize, usize), ArrowError> { - let arrow_data_len = encoded.arrow_data.len(); - if arrow_data_len % usize::from(write_options.alignment) != 0 { - return Err(ArrowError::MemoryError( - "Arrow data not aligned".to_string(), - )); - } - - let a = usize::from(write_options.alignment - 1); - let buffer = encoded.ipc_message; - let flatbuf_size = buffer.len(); - let prefix_size = if write_options.write_legacy_ipc_format { - 4 - } else { - 8 - }; - let aligned_size = (flatbuf_size + prefix_size + a) & !a; - let padding_bytes = aligned_size - flatbuf_size - prefix_size; - - write_continuation( - &mut writer, - write_options, - (aligned_size - prefix_size) as i32, - )?; - - // write the flatbuf - if flatbuf_size > 0 { - writer.write_all(&buffer)?; - } - // write padding - writer.write_all(&PADDING[..padding_bytes])?; - - // write arrow data - let body_len = if arrow_data_len > 0 { - write_body_buffers(&mut writer, &encoded.arrow_data, write_options.alignment)? - } else { - 0 - }; - - Ok((aligned_size, body_len)) -} - -fn write_body_buffers( - mut writer: W, - data: &[u8], - alignment: u8, -) -> Result { - let len = data.len(); - let pad_len = pad_to_alignment(alignment, len); - let total_len = len + pad_len; - - // write body buffer - writer.write_all(data)?; - if pad_len > 0 { - writer.write_all(&PADDING[..pad_len])?; - } - - Ok(total_len) -} - -/// Write a record batch to the writer, writing the message size before the message -/// if the record batch is being written to a stream -fn write_continuation( - mut writer: W, - write_options: &IpcWriteOptions, - total_len: i32, -) -> Result { - let mut written = 8; - - // the version of the writer determines whether continuation markers should be added - match write_options.metadata_version { - crate::MetadataVersion::V1 | crate::MetadataVersion::V2 | crate::MetadataVersion::V3 => { - unreachable!("Options with the metadata version cannot be created") - } - crate::MetadataVersion::V4 => { - if !write_options.write_legacy_ipc_format { - // v0.15.0 format - writer.write_all(&CONTINUATION_MARKER)?; - written = 4; - } - writer.write_all(&total_len.to_le_bytes()[..])?; - } - crate::MetadataVersion::V5 => { - // write continuation marker and message length - writer.write_all(&CONTINUATION_MARKER)?; - writer.write_all(&total_len.to_le_bytes()[..])?; - } - z => panic!("Unsupported crate::MetadataVersion {z:?}"), - }; - - Ok(written) + let mut sink = IpcMessageSink::Writer(&mut writer); + sink.write_encoded_data(encoded, write_options) } /// In V4, null types have no validity bitmap @@ -2772,6 +2754,11 @@ mod tests { stream_reader.next().unwrap().unwrap() } + /// Encodes record batches with [`StreamEncoder`] into one contiguous byte vector. + /// + /// This mirrors callers that need IPC stream bytes without a synchronous + /// [`Write`] implementation, such as encoding buffers before forwarding them + /// to an async sink. fn encode_stream( schema: &Schema, batches: &[RecordBatch], From 2397761538fe5a3de2dc13ccd849e823a60ca585 Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Sat, 18 Jul 2026 10:10:31 +0800 Subject: [PATCH 4/8] docs(ipc): add document comments over sink helpers Signed-off-by: Jiawei Zhao --- arrow-ipc/src/writer.rs | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/arrow-ipc/src/writer.rs b/arrow-ipc/src/writer.rs index 4e2ad4e5c7a6..d711697cb168 100644 --- a/arrow-ipc/src/writer.rs +++ b/arrow-ipc/src/writer.rs @@ -219,6 +219,10 @@ impl IpcMessageSink<'_> { Ok(len + pad_len) } + /// Writes an already encoded IPC message with optional contiguous body data. + /// + /// This is used for schema and dictionary messages represented by [`EncodedData`]. + /// Returns the padded metadata length and body length written. fn write_encoded_data( &mut self, encoded: EncodedData, @@ -256,6 +260,11 @@ impl IpcMessageSink<'_> { Ok((padded_header_len, body_len)) } + /// Writes a record batch message from its encoded metadata and body buffers. + /// + /// The body buffers are already materialized as [`EncodedBuffer`] segments, + /// allowing buffer output to preserve uncompressed Arrow buffers. + /// Returns the padded metadata length and body length written. fn write_record_batch( &mut self, metadata: Vec, @@ -288,6 +297,7 @@ impl IpcMessageSink<'_> { Ok((padded_header_len, body_len)) } + /// Writes the IPC end-of-stream marker. fn write_eos(&mut self, write_options: &IpcWriteOptions) -> Result<(), ArrowError> { self.write_continuation(write_options, 0)?; Ok(()) From 8c53a27090c7307570698f80f349f278e0754ad7 Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Fri, 24 Jul 2026 12:21:33 +0800 Subject: [PATCH 5/8] perf(ipc): specialize message sinks Avoid erasing W: Write on the writer path while preserving shared IPC framing for buffer output. Refs #7812 Signed-off-by: Jiawei Zhao --- arrow-ipc/src/writer.rs | 244 ++++++++++++++++++++++++++++------------ 1 file changed, 172 insertions(+), 72 deletions(-) diff --git a/arrow-ipc/src/writer.rs b/arrow-ipc/src/writer.rs index d711697cb168..a0a8022ec171 100644 --- a/arrow-ipc/src/writer.rs +++ b/arrow-ipc/src/writer.rs @@ -139,77 +139,19 @@ impl<'a> IpcBodySink<'a> { /// /// This emits the stream/file framing around the serialized FlatBuffer /// [`crate::Message`] metadata plus its optional body buffers. -enum IpcMessageSink<'a> { - /// Write bytes directly to a synchronous writer. - Writer(&'a mut dyn Write), - /// Accumulate ordered buffers for deferred writing. - Buffers(&'a mut Vec), -} +trait IpcMessageSink { + fn write_vec(&mut self, bytes: Vec) -> Result<(), ArrowError>; -impl IpcMessageSink<'_> { - fn write_vec(&mut self, bytes: Vec) -> Result<(), ArrowError> { - if bytes.is_empty() { - return Ok(()); - } + fn write_encoded_buffer(&mut self, buffer: EncodedBuffer) -> Result<(), ArrowError>; - match self { - Self::Writer(writer) => writer.write_all(&bytes)?, - Self::Buffers(out) => out.push(Buffer::from(bytes)), - } - Ok(()) - } - - fn write_encoded_buffer(&mut self, buffer: EncodedBuffer) -> Result<(), ArrowError> { - match (self, buffer) { - (Self::Writer(writer), buffer) => writer.write_all(buffer.as_slice())?, - (Self::Buffers(out), EncodedBuffer::Raw(buffer)) => out.push(buffer), - (Self::Buffers(out), EncodedBuffer::Compressed(bytes)) => out.push(Buffer::from(bytes)), - } - Ok(()) - } - - fn write_padding(&mut self, len: usize) -> Result<(), ArrowError> { - if len == 0 { - return Ok(()); - } - - match self { - Self::Writer(writer) => writer.write_all(&PADDING[..len])?, - Self::Buffers(out) => out.push(Buffer::from(&PADDING[..len])), - } - Ok(()) - } + fn write_padding(&mut self, len: usize) -> Result<(), ArrowError>; /// Writes the IPC continuation marker and metadata length prefix. fn write_continuation( &mut self, write_options: &IpcWriteOptions, metadata_len: i32, - ) -> Result<(), ArrowError> { - let mut buffer = Vec::with_capacity(8); - - match write_options.metadata_version { - crate::MetadataVersion::V1 - | crate::MetadataVersion::V2 - | crate::MetadataVersion::V3 => { - unreachable!("Options with the metadata version cannot be created") - } - crate::MetadataVersion::V4 => { - if !write_options.write_legacy_ipc_format { - // v0.15.0 format - buffer.extend_from_slice(&CONTINUATION_MARKER); - } - buffer.extend_from_slice(&metadata_len.to_le_bytes()[..]); - } - crate::MetadataVersion::V5 => { - buffer.extend_from_slice(&CONTINUATION_MARKER); - buffer.extend_from_slice(&metadata_len.to_le_bytes()[..]); - } - z => panic!("Unsupported crate::MetadataVersion {z:?}"), - } - - self.write_vec(buffer) - } + ) -> Result<(), ArrowError>; fn write_body_data(&mut self, data: Vec, alignment: u8) -> Result { let len = data.len(); @@ -304,6 +246,158 @@ impl IpcMessageSink<'_> { } } +/// Writes complete framed IPC messages to a synchronous writer. +struct Writer<'a, W: Write> { + writer: &'a mut W, +} + +impl IpcMessageSink for Writer<'_, W> { + fn write_vec(&mut self, bytes: Vec) -> Result<(), ArrowError> { + if !bytes.is_empty() { + self.writer.write_all(&bytes)?; + } + Ok(()) + } + + fn write_encoded_buffer(&mut self, buffer: EncodedBuffer) -> Result<(), ArrowError> { + self.writer.write_all(buffer.as_slice())?; + Ok(()) + } + + fn write_padding(&mut self, len: usize) -> Result<(), ArrowError> { + if len > 0 { + self.writer.write_all(&PADDING[..len])?; + } + Ok(()) + } + + fn write_continuation( + &mut self, + write_options: &IpcWriteOptions, + metadata_len: i32, + ) -> Result<(), ArrowError> { + let mut buffer = [0; 8]; + let len = match write_options.metadata_version { + crate::MetadataVersion::V1 + | crate::MetadataVersion::V2 + | crate::MetadataVersion::V3 => { + unreachable!("Options with the metadata version cannot be created") + } + crate::MetadataVersion::V4 => { + let metadata_len_bytes = metadata_len.to_le_bytes(); + if !write_options.write_legacy_ipc_format { + // v0.15.0 format + buffer[..4].copy_from_slice(&CONTINUATION_MARKER); + buffer[4..].copy_from_slice(&metadata_len_bytes); + 8 + } else { + buffer[..4].copy_from_slice(&metadata_len_bytes); + 4 + } + } + crate::MetadataVersion::V5 => { + buffer[..4].copy_from_slice(&CONTINUATION_MARKER); + buffer[4..].copy_from_slice(&metadata_len.to_le_bytes()); + 8 + } + z => panic!("Unsupported crate::MetadataVersion {z:?}"), + }; + self.writer.write_all(&buffer[..len])?; + Ok(()) + } + + fn write_record_batch( + &mut self, + metadata: Vec, + encoded_buffers: Vec, + body_len: usize, + tail_pad: usize, + write_options: &IpcWriteOptions, + ) -> Result<(usize, usize), ArrowError> { + let alignment = write_options.alignment; + let prefix_size = if write_options.write_legacy_ipc_format { + 4 + } else { + 8 + }; + let alignment_mask = usize::from(alignment - 1); + let padded_header_len = (metadata.len() + prefix_size + alignment_mask) & !alignment_mask; + let padded_metadata_len = padded_header_len - prefix_size; + let metadata_padding = padded_metadata_len - metadata.len(); + + self.write_continuation(write_options, padded_metadata_len as i32)?; + self.writer.write_all(&metadata)?; + self.writer.write_all(&PADDING[..metadata_padding])?; + for enc in &encoded_buffers { + self.writer.write_all(enc.as_slice())?; + self.writer + .write_all(&PADDING[..pad_to_alignment(alignment, enc.len())])?; + } + self.writer.write_all(&PADDING[..tail_pad])?; + + Ok((padded_header_len, body_len)) + } +} + +/// Accumulates complete framed IPC messages as ordered buffers. +struct Buffers<'a> { + out: &'a mut Vec, +} + +impl IpcMessageSink for Buffers<'_> { + fn write_vec(&mut self, bytes: Vec) -> Result<(), ArrowError> { + if !bytes.is_empty() { + self.out.push(Buffer::from(bytes)); + } + Ok(()) + } + + fn write_encoded_buffer(&mut self, buffer: EncodedBuffer) -> Result<(), ArrowError> { + match buffer { + EncodedBuffer::Raw(buffer) => self.out.push(buffer), + EncodedBuffer::Compressed(bytes) => self.out.push(Buffer::from(bytes)), + } + Ok(()) + } + + fn write_padding(&mut self, len: usize) -> Result<(), ArrowError> { + if len > 0 { + self.out.push(Buffer::from(&PADDING[..len])); + } + Ok(()) + } + + fn write_continuation( + &mut self, + write_options: &IpcWriteOptions, + metadata_len: i32, + ) -> Result<(), ArrowError> { + let mut buffer = Vec::with_capacity(8); + + match write_options.metadata_version { + crate::MetadataVersion::V1 + | crate::MetadataVersion::V2 + | crate::MetadataVersion::V3 => { + unreachable!("Options with the metadata version cannot be created") + } + crate::MetadataVersion::V4 => { + if !write_options.write_legacy_ipc_format { + // v0.15.0 format + buffer.extend_from_slice(&CONTINUATION_MARKER); + } + buffer.extend_from_slice(&metadata_len.to_le_bytes()[..]); + } + crate::MetadataVersion::V5 => { + buffer.extend_from_slice(&CONTINUATION_MARKER); + buffer.extend_from_slice(&metadata_len.to_le_bytes()[..]); + } + z => panic!("Unsupported crate::MetadataVersion {z:?}"), + } + + self.write_vec(buffer) + } +} + /// Per-message sizes produced by [`IpcDataGenerator::write`]. /// /// [`FileWriter`] uses these to build the Block index entries required by the IPC footer for @@ -845,7 +939,7 @@ impl IpcDataGenerator { ipc_write_context: &mut IpcWriteContext, writer: &mut W, ) -> Result { - let mut sink = IpcMessageSink::Writer(writer); + let mut sink = Writer { writer }; self.write_to_sink( batch, dictionary_tracker, @@ -865,7 +959,7 @@ impl IpcDataGenerator { ipc_write_context: &mut IpcWriteContext, out: &mut Vec, ) -> Result { - let mut sink = IpcMessageSink::Buffers(out); + let mut sink = Buffers { out }; self.write_to_sink( batch, dictionary_tracker, @@ -875,13 +969,13 @@ impl IpcDataGenerator { ) } - fn write_to_sink( + fn write_to_sink( &self, batch: &RecordBatch, dictionary_tracker: &mut DictionaryTracker, write_options: &IpcWriteOptions, ipc_write_context: &mut IpcWriteContext, - sink: &mut IpcMessageSink<'_>, + sink: &mut S, ) -> Result { let encoded_dictionaries = self.encode_all_dicts(batch, dictionary_tracker, write_options, ipc_write_context)?; @@ -1665,7 +1759,9 @@ impl FileWriter { // write EOS { - let mut sink = IpcMessageSink::Writer(&mut self.writer); + let mut sink = Writer { + writer: &mut self.writer, + }; sink.write_eos(&self.write_options)?; } @@ -1868,7 +1964,7 @@ impl StreamEncoder { let mut out = vec![]; self.encode_schema(&mut out)?; - let mut sink = IpcMessageSink::Buffers(&mut out); + let mut sink = Buffers { out: &mut out }; sink.write_eos(&self.write_options)?; self.finished = true; Ok(out) @@ -1881,7 +1977,7 @@ impl StreamEncoder { &mut self.dictionary_tracker, &self.write_options, ); - let mut sink = IpcMessageSink::Buffers(out); + let mut sink = Buffers { out }; sink.write_encoded_data(encoded_message, &self.write_options)?; self.schema_encoded = true; } @@ -2058,7 +2154,9 @@ impl StreamWriter { } { - let mut sink = IpcMessageSink::Writer(&mut self.writer); + let mut sink = Writer { + writer: &mut self.writer, + }; sink.write_eos(&self.write_options)?; } self.writer.flush()?; @@ -2158,7 +2256,9 @@ pub fn write_message( encoded: EncodedData, write_options: &IpcWriteOptions, ) -> Result<(usize, usize), ArrowError> { - let mut sink = IpcMessageSink::Writer(&mut writer); + let mut sink = Writer { + writer: &mut writer, + }; sink.write_encoded_data(encoded, write_options) } From 8232385d3cf8fd2859e58055f4a312746ed86492 Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Fri, 24 Jul 2026 15:28:37 +0800 Subject: [PATCH 6/8] refactor(ipc): simplify message sink Implement the sink trait for Write directly so the writer path no longer needs a wrapper sink. Consume StreamEncoder in finish to remove the closed-state check. Refs #7812 Signed-off-by: Jiawei Zhao --- arrow-ipc/src/writer.rs | 202 +++++++++++++--------------------------- 1 file changed, 66 insertions(+), 136 deletions(-) diff --git a/arrow-ipc/src/writer.rs b/arrow-ipc/src/writer.rs index a0a8022ec171..d47295b94e38 100644 --- a/arrow-ipc/src/writer.rs +++ b/arrow-ipc/src/writer.rs @@ -140,18 +140,54 @@ impl<'a> IpcBodySink<'a> { /// This emits the stream/file framing around the serialized FlatBuffer /// [`crate::Message`] metadata plus its optional body buffers. trait IpcMessageSink { - fn write_vec(&mut self, bytes: Vec) -> Result<(), ArrowError>; + fn write_slice(&mut self, bytes: &[u8]) -> Result<(), ArrowError>; - fn write_encoded_buffer(&mut self, buffer: EncodedBuffer) -> Result<(), ArrowError>; + fn write_vec(&mut self, bytes: Vec) -> Result<(), ArrowError> { + self.write_slice(&bytes) + } - fn write_padding(&mut self, len: usize) -> Result<(), ArrowError>; + fn write_encoded_buffer(&mut self, buffer: EncodedBuffer) -> Result<(), ArrowError> { + self.write_slice(buffer.as_slice()) + } + + fn write_padding(&mut self, len: usize) -> Result<(), ArrowError> { + self.write_slice(&PADDING[..len]) + } /// Writes the IPC continuation marker and metadata length prefix. fn write_continuation( &mut self, write_options: &IpcWriteOptions, metadata_len: i32, - ) -> Result<(), ArrowError>; + ) -> Result<(), ArrowError> { + let mut buffer = [0; 8]; + let len = match write_options.metadata_version { + crate::MetadataVersion::V1 + | crate::MetadataVersion::V2 + | crate::MetadataVersion::V3 => { + unreachable!("Options with the metadata version cannot be created") + } + crate::MetadataVersion::V4 => { + let metadata_len_bytes = metadata_len.to_le_bytes(); + if !write_options.write_legacy_ipc_format { + // v0.15.0 format + buffer[..4].copy_from_slice(&CONTINUATION_MARKER); + buffer[4..].copy_from_slice(&metadata_len_bytes); + 8 + } else { + buffer[..4].copy_from_slice(&metadata_len_bytes); + 4 + } + } + crate::MetadataVersion::V5 => { + buffer[..4].copy_from_slice(&CONTINUATION_MARKER); + buffer[4..].copy_from_slice(&metadata_len.to_le_bytes()); + 8 + } + z => panic!("Unsupported crate::MetadataVersion {z:?}"), + }; + self.write_slice(&buffer[..len]) + } fn write_body_data(&mut self, data: Vec, alignment: u8) -> Result { let len = data.len(); @@ -246,66 +282,17 @@ trait IpcMessageSink { } } -/// Writes complete framed IPC messages to a synchronous writer. -struct Writer<'a, W: Write> { - writer: &'a mut W, -} - -impl IpcMessageSink for Writer<'_, W> { - fn write_vec(&mut self, bytes: Vec) -> Result<(), ArrowError> { +impl IpcMessageSink for W +where + W: Write, +{ + fn write_slice(&mut self, bytes: &[u8]) -> Result<(), ArrowError> { if !bytes.is_empty() { - self.writer.write_all(&bytes)?; + self.write_all(bytes)?; } Ok(()) } - fn write_encoded_buffer(&mut self, buffer: EncodedBuffer) -> Result<(), ArrowError> { - self.writer.write_all(buffer.as_slice())?; - Ok(()) - } - - fn write_padding(&mut self, len: usize) -> Result<(), ArrowError> { - if len > 0 { - self.writer.write_all(&PADDING[..len])?; - } - Ok(()) - } - - fn write_continuation( - &mut self, - write_options: &IpcWriteOptions, - metadata_len: i32, - ) -> Result<(), ArrowError> { - let mut buffer = [0; 8]; - let len = match write_options.metadata_version { - crate::MetadataVersion::V1 - | crate::MetadataVersion::V2 - | crate::MetadataVersion::V3 => { - unreachable!("Options with the metadata version cannot be created") - } - crate::MetadataVersion::V4 => { - let metadata_len_bytes = metadata_len.to_le_bytes(); - if !write_options.write_legacy_ipc_format { - // v0.15.0 format - buffer[..4].copy_from_slice(&CONTINUATION_MARKER); - buffer[4..].copy_from_slice(&metadata_len_bytes); - 8 - } else { - buffer[..4].copy_from_slice(&metadata_len_bytes); - 4 - } - } - crate::MetadataVersion::V5 => { - buffer[..4].copy_from_slice(&CONTINUATION_MARKER); - buffer[4..].copy_from_slice(&metadata_len.to_le_bytes()); - 8 - } - z => panic!("Unsupported crate::MetadataVersion {z:?}"), - }; - self.writer.write_all(&buffer[..len])?; - Ok(()) - } - fn write_record_batch( &mut self, metadata: Vec, @@ -326,14 +313,13 @@ impl IpcMessageSink for Writer<'_, W> { let metadata_padding = padded_metadata_len - metadata.len(); self.write_continuation(write_options, padded_metadata_len as i32)?; - self.writer.write_all(&metadata)?; - self.writer.write_all(&PADDING[..metadata_padding])?; + self.write_all(&metadata)?; + self.write_all(&PADDING[..metadata_padding])?; for enc in &encoded_buffers { - self.writer.write_all(enc.as_slice())?; - self.writer - .write_all(&PADDING[..pad_to_alignment(alignment, enc.len())])?; + self.write_all(enc.as_slice())?; + self.write_all(&PADDING[..pad_to_alignment(alignment, enc.len())])?; } - self.writer.write_all(&PADDING[..tail_pad])?; + self.write_all(&PADDING[..tail_pad])?; Ok((padded_header_len, body_len)) } @@ -345,6 +331,13 @@ struct Buffers<'a> { } impl IpcMessageSink for Buffers<'_> { + fn write_slice(&mut self, bytes: &[u8]) -> Result<(), ArrowError> { + if !bytes.is_empty() { + self.out.push(Buffer::from(bytes)); + } + Ok(()) + } + fn write_vec(&mut self, bytes: Vec) -> Result<(), ArrowError> { if !bytes.is_empty() { self.out.push(Buffer::from(bytes)); @@ -359,43 +352,6 @@ impl IpcMessageSink for Buffers<'_> { } Ok(()) } - - fn write_padding(&mut self, len: usize) -> Result<(), ArrowError> { - if len > 0 { - self.out.push(Buffer::from(&PADDING[..len])); - } - Ok(()) - } - - fn write_continuation( - &mut self, - write_options: &IpcWriteOptions, - metadata_len: i32, - ) -> Result<(), ArrowError> { - let mut buffer = Vec::with_capacity(8); - - match write_options.metadata_version { - crate::MetadataVersion::V1 - | crate::MetadataVersion::V2 - | crate::MetadataVersion::V3 => { - unreachable!("Options with the metadata version cannot be created") - } - crate::MetadataVersion::V4 => { - if !write_options.write_legacy_ipc_format { - // v0.15.0 format - buffer.extend_from_slice(&CONTINUATION_MARKER); - } - buffer.extend_from_slice(&metadata_len.to_le_bytes()[..]); - } - crate::MetadataVersion::V5 => { - buffer.extend_from_slice(&CONTINUATION_MARKER); - buffer.extend_from_slice(&metadata_len.to_le_bytes()[..]); - } - z => panic!("Unsupported crate::MetadataVersion {z:?}"), - } - - self.write_vec(buffer) - } } /// Per-message sizes produced by [`IpcDataGenerator::write`]. @@ -939,13 +895,12 @@ impl IpcDataGenerator { ipc_write_context: &mut IpcWriteContext, writer: &mut W, ) -> Result { - let mut sink = Writer { writer }; self.write_to_sink( batch, dictionary_tracker, write_options, ipc_write_context, - &mut sink, + writer, ) } @@ -1759,10 +1714,7 @@ impl FileWriter { // write EOS { - let mut sink = Writer { - writer: &mut self.writer, - }; - sink.write_eos(&self.write_options)?; + self.writer.write_eos(&self.write_options)?; } let mut fbb = FlatBufferBuilder::new(); @@ -1886,8 +1838,6 @@ pub struct StreamEncoder { write_options: IpcWriteOptions, /// Whether the stream schema has been encoded schema_encoded: bool, - /// Whether the end-of-stream marker has been encoded - finished: bool, /// Keeps track of dictionaries that have been encoded dictionary_tracker: DictionaryTracker, data_gen: IpcDataGenerator, @@ -1912,7 +1862,6 @@ impl StreamEncoder { schema: schema.clone(), write_options, schema_encoded: false, - finished: false, dictionary_tracker: DictionaryTracker::new(false), data_gen: IpcDataGenerator::default(), ipc_write_context: IpcWriteContext::default(), @@ -1927,14 +1876,8 @@ impl StreamEncoder { /// /// # Errors /// - /// Returns an error if the encoder is already finished or encoding fails. + /// Returns an error if encoding fails. pub fn encode(&mut self, batch: &RecordBatch) -> Result, ArrowError> { - if self.finished { - return Err(ArrowError::IpcError( - "Cannot encode record batch to stream encoder as it is closed".to_string(), - )); - } - let mut out = vec![]; self.encode_schema(&mut out)?; self.data_gen.encode_to_buffers( @@ -1947,26 +1890,19 @@ impl StreamEncoder { Ok(out) } - /// Encode the end-of-stream marker and mark this encoder as finished. + /// Encode the end-of-stream marker. /// /// If no batches have been encoded, this also emits the IPC stream schema /// message so the returned buffers form a valid empty IPC stream. /// /// # Errors /// - /// Returns an error if the encoder is already finished. - pub fn finish(&mut self) -> Result, ArrowError> { - if self.finished { - return Err(ArrowError::IpcError( - "Cannot finish stream encoder as it is closed".to_string(), - )); - } - + /// Returns an error if encoding the schema or end-of-stream marker fails. + pub fn finish(mut self) -> Result, ArrowError> { let mut out = vec![]; self.encode_schema(&mut out)?; let mut sink = Buffers { out: &mut out }; sink.write_eos(&self.write_options)?; - self.finished = true; Ok(out) } @@ -2154,10 +2090,7 @@ impl StreamWriter { } { - let mut sink = Writer { - writer: &mut self.writer, - }; - sink.write_eos(&self.write_options)?; + self.writer.write_eos(&self.write_options)?; } self.writer.flush()?; @@ -2256,10 +2189,7 @@ pub fn write_message( encoded: EncodedData, write_options: &IpcWriteOptions, ) -> Result<(usize, usize), ArrowError> { - let mut sink = Writer { - writer: &mut writer, - }; - sink.write_encoded_data(encoded, write_options) + writer.write_encoded_data(encoded, write_options) } /// In V4, null types have no validity bitmap From 45e054339456e4c4d19856769ab0043e3f8314e1 Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Fri, 24 Jul 2026 20:37:17 +0800 Subject: [PATCH 7/8] refactor(ipc): share metadata layout Extract the IPC metadata padding calculation so writer and buffer sinks use the same framing layout. Signed-off-by: Jiawei Zhao --- arrow-ipc/src/writer.rs | 73 +++++++++++++++++++++-------------------- 1 file changed, 37 insertions(+), 36 deletions(-) diff --git a/arrow-ipc/src/writer.rs b/arrow-ipc/src/writer.rs index d47295b94e38..a3b6557067a0 100644 --- a/arrow-ipc/src/writer.rs +++ b/arrow-ipc/src/writer.rs @@ -135,6 +135,31 @@ impl<'a> IpcBodySink<'a> { } } +struct MetadataLayout { + padded_header_len: usize, + padded_metadata_len: usize, + metadata_padding: usize, +} + +#[inline] +fn metadata_layout(metadata_len: usize, write_options: &IpcWriteOptions) -> MetadataLayout { + let prefix_size = if write_options.write_legacy_ipc_format { + 4 + } else { + 8 + }; + let alignment_mask = usize::from(write_options.alignment - 1); + let padded_header_len = (metadata_len + prefix_size + alignment_mask) & !alignment_mask; + let padded_metadata_len = padded_header_len - prefix_size; + let metadata_padding = padded_metadata_len - metadata_len; + + MetadataLayout { + padded_header_len, + padded_metadata_len, + metadata_padding, + } +} + /// Destination for a complete framed IPC message. /// /// This emits the stream/file framing around the serialized FlatBuffer @@ -213,21 +238,13 @@ trait IpcMessageSink { )); } - let alignment_mask = usize::from(write_options.alignment - 1); let metadata = encoded.ipc_message; let metadata_len = metadata.len(); - let prefix_size = if write_options.write_legacy_ipc_format { - 4 - } else { - 8 - }; - let padded_header_len = (metadata_len + prefix_size + alignment_mask) & !alignment_mask; - let padded_metadata_len = padded_header_len - prefix_size; - let metadata_padding = padded_metadata_len - metadata_len; + let layout = metadata_layout(metadata_len, write_options); - self.write_continuation(write_options, padded_metadata_len as i32)?; + self.write_continuation(write_options, layout.padded_metadata_len as i32)?; self.write_vec(metadata)?; - self.write_padding(metadata_padding)?; + self.write_padding(layout.metadata_padding)?; let body_len = if arrow_data_len > 0 { self.write_body_data(encoded.arrow_data, write_options.alignment)? @@ -235,7 +252,7 @@ trait IpcMessageSink { 0 }; - Ok((padded_header_len, body_len)) + Ok((layout.padded_header_len, body_len)) } /// Writes a record batch message from its encoded metadata and body buffers. @@ -252,19 +269,11 @@ trait IpcMessageSink { write_options: &IpcWriteOptions, ) -> Result<(usize, usize), ArrowError> { let alignment = write_options.alignment; - let prefix_size = if write_options.write_legacy_ipc_format { - 4 - } else { - 8 - }; - let alignment_mask = usize::from(alignment - 1); - let padded_header_len = (metadata.len() + prefix_size + alignment_mask) & !alignment_mask; - let padded_metadata_len = padded_header_len - prefix_size; - let metadata_padding = padded_metadata_len - metadata.len(); + let layout = metadata_layout(metadata.len(), write_options); - self.write_continuation(write_options, padded_metadata_len as i32)?; + self.write_continuation(write_options, layout.padded_metadata_len as i32)?; self.write_vec(metadata)?; - self.write_padding(metadata_padding)?; + self.write_padding(layout.metadata_padding)?; for enc in encoded_buffers { let len = enc.len(); self.write_encoded_buffer(enc)?; @@ -272,7 +281,7 @@ trait IpcMessageSink { } self.write_padding(tail_pad)?; - Ok((padded_header_len, body_len)) + Ok((layout.padded_header_len, body_len)) } /// Writes the IPC end-of-stream marker. @@ -302,26 +311,18 @@ where write_options: &IpcWriteOptions, ) -> Result<(usize, usize), ArrowError> { let alignment = write_options.alignment; - let prefix_size = if write_options.write_legacy_ipc_format { - 4 - } else { - 8 - }; - let alignment_mask = usize::from(alignment - 1); - let padded_header_len = (metadata.len() + prefix_size + alignment_mask) & !alignment_mask; - let padded_metadata_len = padded_header_len - prefix_size; - let metadata_padding = padded_metadata_len - metadata.len(); + let layout = metadata_layout(metadata.len(), write_options); - self.write_continuation(write_options, padded_metadata_len as i32)?; + self.write_continuation(write_options, layout.padded_metadata_len as i32)?; self.write_all(&metadata)?; - self.write_all(&PADDING[..metadata_padding])?; + self.write_all(&PADDING[..layout.metadata_padding])?; for enc in &encoded_buffers { self.write_all(enc.as_slice())?; self.write_all(&PADDING[..pad_to_alignment(alignment, enc.len())])?; } self.write_all(&PADDING[..tail_pad])?; - Ok((padded_header_len, body_len)) + Ok((layout.padded_header_len, body_len)) } } From ee44c9cc9e55b3d569c2b51967060c8933e2d76e Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Sun, 26 Jul 2026 11:22:19 +0800 Subject: [PATCH 8/8] refactor(ipc): split message sink traits Separate sink-specific write primitives from shared IPC framing helpers while preserving the record batch writer specialization. Signed-off-by: Jiawei Zhao --- arrow-ipc/src/writer.rs | 82 ++++++++++++++++++++++++----------------- 1 file changed, 49 insertions(+), 33 deletions(-) diff --git a/arrow-ipc/src/writer.rs b/arrow-ipc/src/writer.rs index a3b6557067a0..7cdb4618f8eb 100644 --- a/arrow-ipc/src/writer.rs +++ b/arrow-ipc/src/writer.rs @@ -141,29 +141,31 @@ struct MetadataLayout { metadata_padding: usize, } -#[inline] -fn metadata_layout(metadata_len: usize, write_options: &IpcWriteOptions) -> MetadataLayout { - let prefix_size = if write_options.write_legacy_ipc_format { - 4 - } else { - 8 - }; - let alignment_mask = usize::from(write_options.alignment - 1); - let padded_header_len = (metadata_len + prefix_size + alignment_mask) & !alignment_mask; - let padded_metadata_len = padded_header_len - prefix_size; - let metadata_padding = padded_metadata_len - metadata_len; +impl MetadataLayout { + fn new(metadata_len: usize, write_options: &IpcWriteOptions) -> MetadataLayout { + let prefix_size = if write_options.write_legacy_ipc_format { + 4 + } else { + 8 + }; + let alignment_mask = usize::from(write_options.alignment - 1); + let padded_header_len = (metadata_len + prefix_size + alignment_mask) & !alignment_mask; + let padded_metadata_len = padded_header_len - prefix_size; + let metadata_padding = padded_metadata_len - metadata_len; - MetadataLayout { - padded_header_len, - padded_metadata_len, - metadata_padding, + MetadataLayout { + padded_header_len, + padded_metadata_len, + metadata_padding, + } } } -/// Destination for a complete framed IPC message. +/// Destination-specific byte writing for a framed IPC message. /// -/// This emits the stream/file framing around the serialized FlatBuffer -/// [`crate::Message`] metadata plus its optional body buffers. +/// The default owned-buffer methods borrow their bytes and call `write_slice`, +/// which can copy the data. Sink implementations that can take ownership of +/// generated bytes should override `write_vec` or `write_encoded_buffer`. trait IpcMessageSink { fn write_slice(&mut self, bytes: &[u8]) -> Result<(), ArrowError>; @@ -174,7 +176,10 @@ trait IpcMessageSink { fn write_encoded_buffer(&mut self, buffer: EncodedBuffer) -> Result<(), ArrowError> { self.write_slice(buffer.as_slice()) } +} +/// Shared IPC framing helpers for [`IpcMessageSink`]. +trait IpcMessageSinkExt: IpcMessageSink { fn write_padding(&mut self, len: usize) -> Result<(), ArrowError> { self.write_slice(&PADDING[..len]) } @@ -240,7 +245,7 @@ trait IpcMessageSink { let metadata = encoded.ipc_message; let metadata_len = metadata.len(); - let layout = metadata_layout(metadata_len, write_options); + let layout = MetadataLayout::new(metadata_len, write_options); self.write_continuation(write_options, layout.padded_metadata_len as i32)?; self.write_vec(metadata)?; @@ -255,6 +260,17 @@ trait IpcMessageSink { Ok((layout.padded_header_len, body_len)) } + /// Writes the IPC end-of-stream marker. + fn write_eos(&mut self, write_options: &IpcWriteOptions) -> Result<(), ArrowError> { + self.write_continuation(write_options, 0)?; + Ok(()) + } +} + +impl IpcMessageSinkExt for T {} + +/// Optional hot-path hook for record batch messages. +trait IpcRecordBatchSink: IpcMessageSinkExt { /// Writes a record batch message from its encoded metadata and body buffers. /// /// The body buffers are already materialized as [`EncodedBuffer`] segments, @@ -269,7 +285,7 @@ trait IpcMessageSink { write_options: &IpcWriteOptions, ) -> Result<(usize, usize), ArrowError> { let alignment = write_options.alignment; - let layout = metadata_layout(metadata.len(), write_options); + let layout = MetadataLayout::new(metadata.len(), write_options); self.write_continuation(write_options, layout.padded_metadata_len as i32)?; self.write_vec(metadata)?; @@ -283,12 +299,6 @@ trait IpcMessageSink { Ok((layout.padded_header_len, body_len)) } - - /// Writes the IPC end-of-stream marker. - fn write_eos(&mut self, write_options: &IpcWriteOptions) -> Result<(), ArrowError> { - self.write_continuation(write_options, 0)?; - Ok(()) - } } impl IpcMessageSink for W @@ -301,7 +311,12 @@ where } Ok(()) } +} +impl IpcRecordBatchSink for W +where + W: Write, +{ fn write_record_batch( &mut self, metadata: Vec, @@ -311,7 +326,7 @@ where write_options: &IpcWriteOptions, ) -> Result<(usize, usize), ArrowError> { let alignment = write_options.alignment; - let layout = metadata_layout(metadata.len(), write_options); + let layout = MetadataLayout::new(metadata.len(), write_options); self.write_continuation(write_options, layout.padded_metadata_len as i32)?; self.write_all(&metadata)?; @@ -355,6 +370,8 @@ impl IpcMessageSink for Buffers<'_> { } } +impl IpcRecordBatchSink for Buffers<'_> {} + /// Per-message sizes produced by [`IpcDataGenerator::write`]. /// /// [`FileWriter`] uses these to build the Block index entries required by the IPC footer for @@ -925,7 +942,7 @@ impl IpcDataGenerator { ) } - fn write_to_sink( + fn write_to_sink( &self, batch: &RecordBatch, dictionary_tracker: &mut DictionaryTracker, @@ -2797,9 +2814,8 @@ mod tests { /// Encodes record batches with [`StreamEncoder`] into one contiguous byte vector. /// - /// This mirrors callers that need IPC stream bytes without a synchronous - /// [`Write`] implementation, such as encoding buffers before forwarding them - /// to an async sink. + /// This still exercises the sans-IO encoder path; `Vec` is only used + /// as a convenient [`Write`] sink for comparing the resulting byte stream. fn encode_stream( schema: &Schema, batches: &[RecordBatch], @@ -2809,11 +2825,11 @@ mod tests { let mut bytes = Vec::new(); for batch in batches { for buffer in encoder.encode(batch).unwrap() { - bytes.extend_from_slice(buffer.as_slice()); + bytes.write_all(buffer.as_slice()).unwrap(); } } for buffer in encoder.finish().unwrap() { - bytes.extend_from_slice(buffer.as_slice()); + bytes.write_all(buffer.as_slice()).unwrap(); } bytes }