leaves12138 commented on code in PR #959:
URL: https://github.com/apache/paimon-rust/pull/959#discussion_r4109931199


##########
crates/paimon/src/arrow/format/text.rs:
##########
@@ -0,0 +1,1310 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Line-oriented Format Table files. Java uses a positional CSV schema, JSON
+//! objects keyed by field name, and exactly one string column for TEXT.
+
+use super::{FilePredicates, FormatFileReader, FormatFileWriter, 
FormatWriteResult};
+use crate::arrow::build_target_arrow_schema;
+use crate::io::{FileRead, FileWrite, OutputFile};
+use crate::spec::DataField;
+use crate::table::{ArrowRecordBatchStream, RowRange};
+use crate::Error;
+use arrow_array::{
+    Array, ArrayRef, BinaryArray, FixedSizeBinaryArray, LargeBinaryArray, 
RecordBatch, StringArray,
+};
+use arrow_schema::{DataType, SchemaRef};
+use async_trait::async_trait;
+use base64::Engine;
+use bytes::Bytes;
+use futures::{stream, StreamExt};
+use std::collections::HashMap;
+use std::io::{self, BufReader, Cursor, Read, Write};
+use std::sync::Arc;
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(super) enum TextKind {
+    Csv,
+    Json,
+    Text,
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(crate) enum TextCompression {
+    None,
+    Gzip,
+    Bzip2,
+    Deflate,
+    Snappy,
+    Lz4,
+    Zstd,
+}
+
+impl TextCompression {
+    pub(crate) fn from_name(name: &str) -> crate::Result<Self> {
+        match name.to_ascii_lowercase().as_str() {
+            "" | "none" | "uncompressed" => Ok(Self::None),
+            "gzip" => Ok(Self::Gzip),
+            "bzip2" => Ok(Self::Bzip2),
+            "deflate" => Ok(Self::Deflate),
+            "snappy" => Ok(Self::Snappy),
+            "lz4" => Ok(Self::Lz4),
+            "zstd" => Ok(Self::Zstd),
+            _ => Err(Error::Unsupported {
+                message: format!("Unsupported text file compression '{name}'"),
+            }),
+        }
+    }
+
+    pub(crate) fn extension(self) -> Option<&'static str> {
+        match self {
+            Self::None => None,
+            Self::Gzip => Some("gz"),
+            Self::Bzip2 => Some("bz2"),
+            Self::Deflate => Some("deflate"),
+            Self::Snappy => Some("snappy"),
+            Self::Lz4 => Some("lz4"),
+            Self::Zstd => Some("zst"),
+        }
+    }
+
+    fn from_extension(extension: &str) -> Option<Self> {
+        match extension {
+            "gz" => Some(Self::Gzip),
+            "bz2" => Some(Self::Bzip2),
+            "deflate" => Some(Self::Deflate),
+            "snappy" => Some(Self::Snappy),
+            "lz4" => Some(Self::Lz4),
+            "zst" => Some(Self::Zstd),
+            _ => None,
+        }
+    }
+}
+
+enum TextEncoder {
+    None,
+    Gzip(flate2::write::GzEncoder<Vec<u8>>),
+    Bzip2(bzip2::write::BzEncoder<Vec<u8>>),
+    Deflate(flate2::write::ZlibEncoder<Vec<u8>>),
+    Snappy,
+    Lz4,
+    Zstd(zstd::stream::write::Encoder<'static, Vec<u8>>),
+}
+
+impl TextEncoder {
+    fn new(compression: TextCompression) -> crate::Result<Self> {
+        Ok(match compression {
+            TextCompression::None => Self::None,
+            TextCompression::Gzip => Self::Gzip(flate2::write::GzEncoder::new(
+                Vec::new(),
+                flate2::Compression::default(),
+            )),
+            TextCompression::Bzip2 => Self::Bzip2(bzip2::write::BzEncoder::new(
+                Vec::new(),
+                bzip2::Compression::default(),
+            )),
+            TextCompression::Deflate => 
Self::Deflate(flate2::write::ZlibEncoder::new(
+                Vec::new(),
+                flate2::Compression::default(),
+            )),
+            TextCompression::Snappy => Self::Snappy,
+            TextCompression::Lz4 => Self::Lz4,
+            TextCompression::Zstd => Self::Zstd(
+                zstd::stream::write::Encoder::new(Vec::new(), 
0).map_err(compression_error)?,
+            ),
+        })
+    }
+
+    fn write(&mut self, input: Vec<u8>) -> crate::Result<Vec<u8>> {
+        match self {
+            Self::None => Ok(input),
+            Self::Gzip(encoder) => {
+                encoder.write_all(&input).map_err(compression_error)?;
+                Ok(std::mem::take(encoder.get_mut()))
+            }
+            Self::Bzip2(encoder) => {
+                encoder.write_all(&input).map_err(compression_error)?;
+                Ok(std::mem::take(encoder.get_mut()))
+            }
+            Self::Deflate(encoder) => {
+                encoder.write_all(&input).map_err(compression_error)?;
+                Ok(std::mem::take(encoder.get_mut()))
+            }
+            Self::Snappy => encode_hadoop_blocks(&input, 
TextCompression::Snappy),
+            Self::Lz4 => encode_hadoop_blocks(&input, TextCompression::Lz4),
+            Self::Zstd(encoder) => {
+                encoder.write_all(&input).map_err(compression_error)?;
+                Ok(std::mem::take(encoder.get_mut()))
+            }
+        }
+    }
+
+    fn finish(self) -> crate::Result<Vec<u8>> {
+        match self {
+            Self::None => Ok(Vec::new()),
+            Self::Gzip(encoder) => encoder.finish().map_err(compression_error),
+            Self::Bzip2(encoder) => 
encoder.finish().map_err(compression_error),
+            Self::Deflate(encoder) => 
encoder.finish().map_err(compression_error),
+            Self::Snappy | Self::Lz4 => Ok(0_u32.to_be_bytes().to_vec()),
+            Self::Zstd(encoder) => encoder.finish().map_err(compression_error),
+        }
+    }
+}
+
+impl TextKind {
+    fn extension(self) -> &'static str {
+        match self {
+            Self::Csv => ".csv",
+            Self::Json => ".json",
+            Self::Text => ".text",
+        }
+    }
+
+    pub(super) fn from_path(path: &str) -> Option<(Self, TextCompression)> {
+        let path = path.to_ascii_lowercase();
+        for kind in [Self::Csv, Self::Json, Self::Text] {
+            if path.ends_with(kind.extension()) {
+                return Some((kind, TextCompression::None));
+            }
+            if let Some((_, suffix)) = path.rsplit_once(&format!("{}.", 
kind.extension())) {
+                if let Some(codec) = TextCompression::from_extension(suffix) {
+                    return Some((kind, codec));
+                }
+            }
+        }
+        None
+    }
+}
+
+pub(crate) fn matches_compressed_extension(file_name: &str, format_extension: 
&str) -> bool {
+    [".gz", ".bz2", ".deflate", ".snappy", ".lz4", ".zst"]
+        .iter()
+        .any(|suffix| {
+            file_name
+                .strip_suffix(suffix)
+                .is_some_and(|stem| stem.ends_with(format_extension))
+        })
+}
+
+fn compression_error(error: std::io::Error) -> Error {
+    Error::DataInvalid {
+        message: format!("Invalid compressed text file: {error}"),
+        source: Some(Box::new(error)),
+    }
+}
+
+// Hadoop's SnappyCodec and Lz4Codec wrap raw codec blocks with a big-endian
+// uncompressed length followed by one or more length-prefixed compressed 
chunks.
+fn encode_hadoop_blocks(input: &[u8], codec: TextCompression) -> 
crate::Result<Vec<u8>> {
+    const BLOCK_SIZE: usize = 256 * 1024 - 2048;
+    let mut output = Vec::new();
+    for block in input.chunks(BLOCK_SIZE) {
+        let compressed = match codec {
+            TextCompression::Snappy => {
+                snap::raw::Encoder::new()
+                    .compress_vec(block)
+                    .map_err(|error| Error::DataInvalid {
+                        message: format!("Snappy compression failed: {error}"),
+                        source: Some(Box::new(error)),
+                    })?
+            }
+            TextCompression::Lz4 => lz4_flex::block::compress(block),
+            _ => unreachable!(),
+        };
+        output.extend_from_slice(&(block.len() as u32).to_be_bytes());
+        output.extend_from_slice(&(compressed.len() as u32).to_be_bytes());
+        output.extend_from_slice(&compressed);
+    }
+    Ok(output)
+}
+
+/// Decode one Hadoop block at a time. A malformed length must not allocate an
+/// attacker-controlled amount of memory before the file can be rejected.
+struct HadoopBlockReader<R> {
+    reader: R,
+    codec: TextCompression,
+    block: Vec<u8>,
+    position: usize,
+    finished: bool,
+}
+
+impl<R: Read> HadoopBlockReader<R> {
+    fn new(reader: R, codec: TextCompression) -> Self {
+        Self {
+            reader,
+            codec,
+            block: Vec::new(),
+            position: 0,
+            finished: false,
+        }
+    }
+
+    fn load_block(&mut self) -> io::Result<()> {
+        const MAX_BLOCK_SIZE: usize = 64 * 1024 * 1024;
+        let mut length = [0; 4];
+        if self.reader.read(&mut length[..1])? == 0 {
+            self.finished = true;
+            return Ok(());
+        }
+        self.reader.read_exact(&mut length[1..])?;
+        let original_size = u32::from_be_bytes(length) as usize;
+        if original_size == 0 {
+            if self.reader.read(&mut length[..1])? != 0 {
+                return Err(io::Error::new(
+                    io::ErrorKind::InvalidData,
+                    "bytes after final block",
+                ));
+            }
+            self.finished = true;
+            return Ok(());
+        }
+        if original_size > MAX_BLOCK_SIZE {
+            return Err(io::Error::new(
+                io::ErrorKind::InvalidData,
+                "Hadoop block is too large",
+            ));
+        }
+        self.block.clear();
+        while self.block.len() < original_size {
+            self.reader.read_exact(&mut length)?;
+            let compressed_size = u32::from_be_bytes(length) as usize;
+            if compressed_size == 0 || compressed_size > MAX_BLOCK_SIZE {
+                return Err(io::Error::new(
+                    io::ErrorKind::InvalidData,
+                    "invalid compressed chunk length",
+                ));
+            }
+            let mut compressed = vec![0; compressed_size];
+            self.reader.read_exact(&mut compressed)?;
+            let remaining = original_size - self.block.len();
+            let decoded = match self.codec {
+                TextCompression::Snappy => {
+                    let decoded_len = snap::raw::decompress_len(&compressed)
+                        .map_err(|e| 
io::Error::new(io::ErrorKind::InvalidData, e))?;
+                    if decoded_len > remaining {
+                        return Err(io::Error::new(
+                            io::ErrorKind::InvalidData,
+                            "Snappy chunk exceeds block length",
+                        ));
+                    }
+                    snap::raw::Decoder::new()
+                        .decompress_vec(&compressed)
+                        .map_err(|e| 
io::Error::new(io::ErrorKind::InvalidData, e))?
+                }
+                TextCompression::Lz4 => 
lz4_flex::block::decompress(&compressed, remaining)
+                    .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, 
e))?,
+                _ => unreachable!(),
+            };
+            if decoded.is_empty() || decoded.len() > remaining {
+                return Err(io::Error::new(
+                    io::ErrorKind::InvalidData,
+                    "decoded chunk exceeds block length",
+                ));
+            }
+            self.block.extend_from_slice(&decoded);
+        }
+        self.position = 0;
+        Ok(())
+    }
+}
+
+impl<R: Read> Read for HadoopBlockReader<R> {
+    fn read(&mut self, output: &mut [u8]) -> io::Result<usize> {
+        if output.is_empty() {
+            return Ok(0);
+        }
+        if self.position == self.block.len() && !self.finished {
+            self.load_block()?;
+        }
+        if self.finished {
+            return Ok(0);
+        }
+        let size = output.len().min(self.block.len() - self.position);
+        
output[..size].copy_from_slice(&self.block[self.position..self.position + 
size]);
+        self.position += size;
+        Ok(size)
+    }
+}
+
+#[cfg(test)]
+mod compression_tests {
+    use super::*;
+    use crate::spec::{DataType as PaimonDataType, VarCharType};
+    use std::ops::Range;
+    use std::sync::atomic::{AtomicUsize, Ordering};
+
+    struct CountedReader {
+        bytes: Bytes,
+        read_bytes: Arc<AtomicUsize>,
+    }
+
+    #[async_trait]
+    impl FileRead for CountedReader {
+        async fn read(&self, range: Range<u64>) -> crate::Result<Bytes> {
+            self.read_bytes
+                .fetch_add((range.end - range.start) as usize, 
Ordering::SeqCst);
+            Ok(self.bytes.slice(range.start as usize..range.end as usize))
+        }
+    }
+
+    #[test]
+    fn hadoop_block_codecs_read_multiple_blocks_and_reject_truncation() {
+        let input = (0..600_000)
+            .map(|value| (value % 251) as u8)
+            .collect::<Vec<_>>();
+        for codec in [TextCompression::Snappy, TextCompression::Lz4] {
+            let mut encoded = encode_hadoop_blocks(&input, codec).unwrap();
+            encoded.extend_from_slice(&0_u32.to_be_bytes());
+            let mut streamed = Vec::new();
+            HadoopBlockReader::new(Cursor::new(&encoded), codec)
+                .read_to_end(&mut streamed)
+                .unwrap();
+            assert_eq!(streamed, input);
+            encoded.truncate(encoded.len() - 2);
+            assert!(HadoopBlockReader::new(Cursor::new(&encoded), codec)
+                .read_to_end(&mut Vec::new())
+                .is_err());
+        }
+    }
+
+    #[test]
+    fn reads_hadoop_342_java_codec_outputs() {
+        // Produced by Hadoop 3.4.2 SnappyCodec and Lz4Codec from the same 
input.
+        // Hadoop does not add an explicit terminal block in these files.
+        let expected = [b"hello,payload\n".repeat(20), 
(0..16).collect::<Vec<_>>()].concat();
+        for (codec, encoded) in [
+            (
+                TextCompression::Snappy,
+                
"0000012800000030a8023468656c6c6f2c7061796c6f61640afe0e00fe0e00fe0e00fe0e00190e3c000102030405060708090a0b0c0d0e0f",
+            ),
+            (
+                TextCompression::Lz4,
+                
"0000012800000024ef68656c6c6f2c7061796c6f61640a0e00f7f001000102030405060708090a0b0c0d0e0f",
+            ),
+        ] {
+            let encoded = hex::decode(encoded).unwrap();
+            let mut actual = Vec::new();
+            HadoopBlockReader::new(Cursor::new(encoded), codec)
+                .read_to_end(&mut actual)
+                .unwrap();
+            assert_eq!(actual, expected, "{codec:?}");
+        }
+    }
+
+    #[test]
+    fn line_delimiter_crosses_read_chunk_boundary() {
+        let mut input = vec![b'a'; 64 * 1024 - 1];
+        input.extend_from_slice(b"||tail||");
+        let mut lines = DelimitedLines::new(Cursor::new(input), "||");
+        assert_eq!(lines.next_line().unwrap().unwrap().len(), 64 * 1024 - 1);
+        assert_eq!(lines.next_line().unwrap().unwrap(), b"tail");
+        assert!(lines.next_line().unwrap().is_none());
+    }
+
+    #[tokio::test]
+    async fn text_reader_emits_first_batch_before_reading_whole_file() {
+        let bytes = Bytes::from("row\n".repeat(250_000));
+        let size = bytes.len();
+        let read_bytes = Arc::new(AtomicUsize::new(0));
+        let reader =
+            TextFormatReader::new(TextKind::Text, TextCompression::None, 
&HashMap::new()).unwrap();
+        let fields = [DataField::new(
+            0,
+            "line".to_string(),
+            PaimonDataType::VarChar(VarCharType::string_type()),
+        )];
+        let mut stream = reader
+            .read_batch_stream(
+                Box::new(CountedReader {
+                    bytes,
+                    read_bytes: read_bytes.clone(),
+                }),
+                size as u64,
+                &fields,
+                None,
+                Some(1024),
+                None,
+            )
+            .await
+            .unwrap();
+        assert_eq!(stream.next().await.unwrap().unwrap().num_rows(), 1024);
+        assert!(read_bytes.load(Ordering::SeqCst) < size / 2);
+    }
+}
+
+#[derive(Clone)]
+struct TextOptions {
+    line_delimiter: String,
+    field_delimiter: u8,
+    quote: u8,
+    escape: u8,
+    header: bool,
+    null_literal: String,
+}
+
+impl TextOptions {
+    fn new(kind: TextKind, options: &HashMap<String, String>) -> 
crate::Result<Self> {
+        let prefix = match kind {
+            TextKind::Csv => "csv",
+            TextKind::Json => "json",
+            TextKind::Text => "text",
+        };
+        let get = |key: &str, fallback: &str| {
+            options
+                .get(&format!("{prefix}.{key}"))
+                .or_else(|| options.get(fallback))
+                .cloned()
+        };
+        let line_delimiter = get("line-delimiter", 
"lineSep").unwrap_or_else(|| "\n".into());
+        if line_delimiter.is_empty() {
+            return Err(Error::ConfigInvalid {
+                message: format!("{prefix}.line-delimiter must not be empty"),
+            });
+        }
+        let byte = |key: &str, fallback: &str, default: u8| -> 
crate::Result<u8> {
+            let value = get(key, fallback).unwrap_or_else(|| (default as 
char).to_string());
+            if value.len() != 1 {
+                return Err(Error::ConfigInvalid {
+                    message: format!("{prefix}.{key} must be one ASCII 
character"),
+                });
+            }
+            Ok(value.as_bytes()[0])
+        };
+        Ok(Self {
+            line_delimiter,
+            field_delimiter: byte("field-delimiter", "delimiter", b',')?,
+            quote: byte("quote-character", "quote", b'"')?,
+            escape: byte("escape-character", "escape", b'\\')?,
+            header: get("include-header", "header")
+                .is_some_and(|value| value.eq_ignore_ascii_case("true")),
+            null_literal: get("null-literal", "nullvalue").unwrap_or_default(),
+        })
+    }
+}
+
+pub(super) struct TextFormatReader {
+    kind: TextKind,
+    compression: TextCompression,
+    options: TextOptions,
+}
+
+impl TextFormatReader {
+    pub(super) fn new(
+        kind: TextKind,
+        compression: TextCompression,
+        options: &HashMap<String, String>,
+    ) -> crate::Result<Self> {
+        Ok(Self {
+            kind,
+            compression,
+            options: TextOptions::new(kind, options)?,
+        })
+    }
+}
+
+#[async_trait]
+impl FormatFileReader for TextFormatReader {
+    fn select_read_fields(
+        &self,
+        data_schema_fields: &[DataField],
+        projected_fields: &[DataField],
+    ) -> Vec<DataField> {
+        // CSV and TEXT have no field names in their data rows. A projection
+        // still has to decode their original physical column positions.
+        match self.kind {
+            TextKind::Csv | TextKind::Text => data_schema_fields.to_vec(),
+            TextKind::Json => projected_fields.to_vec(),
+        }
+    }
+
+    async fn read_batch_stream(
+        &self,
+        reader: Box<dyn FileRead>,
+        file_size: u64,
+        read_fields: &[DataField],
+        predicates: Option<&FilePredicates>,
+        batch_size: Option<usize>,
+        row_selection: Option<Vec<RowRange>>,
+    ) -> crate::Result<ArrowRecordBatchStream> {
+        if row_selection.is_some() {
+            return Err(Error::Unsupported {
+                message: "Row selection is not supported for line-oriented 
formats".into(),
+            });
+        }
+        let fields = crate::arrow::residual::widen_scan_fields(read_fields, 
predicates);
+        let schema = build_target_arrow_schema(&fields)?;
+        match self.kind {
+            TextKind::Csv => validate_csv_schema(&schema)?,
+            TextKind::Text => validate_text_schema(&schema)?,
+            TextKind::Json => {}
+        }
+        let kind = self.kind;
+        let compression = self.compression;
+        let options = self.options.clone();
+        let size = batch_size.unwrap_or(1024).max(1);
+        let (sender, receiver) = tokio::sync::mpsc::channel(2);
+        let chunks = stream::try_unfold((reader, 0_u64), move |(reader, 
start)| async move {
+            if start >= file_size {
+                return Ok::<_, io::Error>(None);
+            }
+            let end = start.saturating_add(64 * 1024).min(file_size);
+            let bytes = 
reader.read(start..end).await.map_err(io::Error::other)?;
+            if bytes.len() != (end - start) as usize {
+                return Err(io::Error::new(
+                    io::ErrorKind::UnexpectedEof,
+                    "text file changed while reading",
+                ));
+            }
+            Ok(Some((bytes, (reader, end))))
+        });
+        let source =
+            
tokio_util::io::SyncIoBridge::new(tokio_util::io::StreamReader::new(Box::pin(chunks)));
+        tokio::task::spawn_blocking(move || {
+            let source: Box<dyn Read> = match compression {
+                TextCompression::None => Box::new(source),
+                TextCompression::Gzip => 
Box::new(flate2::read::MultiGzDecoder::new(source)),
+                TextCompression::Bzip2 => 
Box::new(bzip2::read::MultiBzDecoder::new(source)),
+                TextCompression::Deflate => 
Box::new(flate2::read::ZlibDecoder::new(source)),
+                TextCompression::Snappy | TextCompression::Lz4 => {
+                    Box::new(HadoopBlockReader::new(source, compression))
+                }
+                TextCompression::Zstd => match 
zstd::stream::read::Decoder::new(source) {
+                    Ok(decoder) => Box::new(decoder),
+                    Err(error) => {
+                        let _ = 
sender.blocking_send(Err(compression_error(error)));
+                        return;
+                    }
+                },
+            };
+            if let Err(error) = decode_text_batches(source, kind, options, 
schema, size, &sender) {
+                let _ = sender.blocking_send(Err(error));
+            }
+        });
+        let predicates = predicates.map(|fp| FilePredicates {
+            predicates: fp.predicates.clone(),
+            row_filter_factory: None,
+            file_fields: fp.file_fields.clone(),
+        });
+        Ok(stream::unfold(receiver, |mut receiver| async {
+            receiver.recv().await.map(|batch| (batch, receiver))
+        })
+        .map(move |batch| {
+            let batch = batch?;
+            match predicates.as_ref() {
+                Some(fp) => {
+                    
crate::arrow::residual::filter_record_batch_by_predicates(batch, fp, &fields)
+                }
+                None => Ok(batch),
+            }
+        })
+        .boxed())
+    }
+}
+
+struct DelimitedLines<R> {
+    reader: BufReader<R>,
+    delimiter: Vec<u8>,
+    pending: Vec<u8>,
+    search_from: usize,
+    eof: bool,
+}
+
+impl<R: Read> DelimitedLines<R> {
+    fn new(reader: R, delimiter: &str) -> Self {
+        Self {
+            reader: BufReader::new(reader),
+            delimiter: delimiter.as_bytes().to_vec(),
+            pending: Vec::new(),
+            search_from: 0,
+            eof: false,
+        }
+    }
+
+    fn next_line(&mut self) -> io::Result<Option<Vec<u8>>> {
+        loop {
+            if let Some(relative) = self.pending[self.search_from..]
+                .windows(self.delimiter.len())
+                .position(|window| window == self.delimiter)
+            {
+                let position = self.search_from + relative;
+                let rest = self.pending.split_off(position + 
self.delimiter.len());
+                let mut line = std::mem::replace(&mut self.pending, rest);
+                line.truncate(position);
+                self.search_from = 0;
+                return Ok(Some(line));
+            }
+            if self.eof {
+                return Ok((!self.pending.is_empty()).then(|| 
std::mem::take(&mut self.pending)));

Review Comment:
   [P1] Reset the search cursor when consuming the final unterminated record
   
   At EOF, `std::mem::take(&mut self.pending)` empties the buffer but leaves 
`search_from` at the end of the old buffer. The next call to `next_line()` 
reaches `self.pending[self.search_from..]` before checking `eof`, and panics 
with an out-of-range slice. This is a normal input case: text files do not have 
to end with a newline.
   
   The impact is silent row loss, not just an error: `read_batch_stream()` 
discards the `spawn_blocking` JoinHandle, so the panic drops the sender and the 
consumer interprets it as clean EOF. Reproduced through the public Format Table 
scan/read API:
   - A one-row CSV or TEXT file containing exactly `hello`, or JSON containing 
exactly `{"value":"hello"}`, returns `Ok` with **0 rows**.
   - A 2,050-row CSV without the final newline returns **2,048 rows**, dropping 
the entire last buffered chunk.
   - A gzip-compressed CSV without the final newline fails the same way.
   
   Adding the final newline makes the controls pass, and the actual Java 
CSV/JSON/TEXT readers correctly return the unterminated record. Please 
reset/check the EOF cursor state before slicing the drained buffer and 
propagate decoder-task failures to the stream rather than treating them as 
successful completion. Add regression coverage for unterminated final records, 
partial final batches, and compressed input.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to