This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new 2e1087a6 fix(vortex): reject columns the writer cannot convert instead
of panicking (#878)
2e1087a6 is described below
commit 2e1087a6537613aa67cbbb5da0a5ef966e6d4cdd
Author: jackylee <[email protected]>
AuthorDate: Tue Sep 22 14:15:44 2026 +0800
fix(vortex): reject columns the writer cannot convert instead of panicking
(#878)
---
crates/paimon/src/arrow/format/vortex.rs | 137 +++++++++++++++++++++++++++++++
1 file changed, 137 insertions(+)
diff --git a/crates/paimon/src/arrow/format/vortex.rs
b/crates/paimon/src/arrow/format/vortex.rs
index 7d6d637e..b53437f3 100644
--- a/crates/paimon/src/arrow/format/vortex.rs
+++ b/crates/paimon/src/arrow/format/vortex.rs
@@ -444,6 +444,7 @@ pub(crate) struct VortexFormatWriter {
impl VortexFormatWriter {
pub(crate) async fn new(output: &OutputFile, schema: SchemaRef) ->
crate::Result<Self> {
+ validate_vortex_schema(&schema)?;
let dtype = DType::from_arrow(schema);
let bytes_written = Arc::new(AtomicU64::new(0));
@@ -457,6 +458,42 @@ impl VortexFormatWriter {
}
}
+/// `DType::from_arrow` ends in `unimplemented!` for a type it has no arm for,
so screen
+/// the schema first. `Map` is the only such type Paimon builds — it is what
+/// `paimon_type_to_arrow` emits for both `MAP` and `MULTISET` — and it can
sit at any
+/// depth, so the walk descends the same containers `DType::from_arrow` does.
+fn validate_vortex_schema(schema: &SchemaRef) -> crate::Result<()> {
+ for field in schema.fields() {
+ if contains_unconvertible_type(field.data_type()) {
+ return Err(Error::Unsupported {
+ message: format!(
+ ".vortex writer does not support column '{}' with type
{:?}",
+ field.name(),
+ field.data_type()
+ ),
+ });
+ }
+ }
+ Ok(())
+}
+
+fn contains_unconvertible_type(data_type: &ArrowDataType) -> bool {
+ match data_type {
+ ArrowDataType::Map(..) => true,
+ ArrowDataType::List(field)
+ | ArrowDataType::LargeList(field)
+ | ArrowDataType::ListView(field)
+ | ArrowDataType::LargeListView(field)
+ | ArrowDataType::FixedSizeList(field, _)
+ | ArrowDataType::RunEndEncoded(_, field) =>
contains_unconvertible_type(field.data_type()),
+ ArrowDataType::Struct(fields) => fields
+ .iter()
+ .any(|field| contains_unconvertible_type(field.data_type())),
+ ArrowDataType::Dictionary(_, value_type) =>
contains_unconvertible_type(value_type),
+ _ => false,
+ }
+}
+
#[async_trait]
impl FormatFileWriter for VortexFormatWriter {
async fn write(&mut self, batch: &RecordBatch) -> crate::Result<()> {
@@ -1392,4 +1429,104 @@ mod tests {
}
total_rows
}
+
+ #[tokio::test]
+ async fn test_vortex_writer_rejects_a_map_column() {
+ use crate::arrow::build_target_arrow_schema;
+ use crate::spec::{DataField, DataType, IntType, MapType, VarCharType};
+
+ let fields = vec![
+ DataField::new(0, "id".to_string(), DataType::Int(IntType::new())),
+ DataField::new(
+ 1,
+ "tags".to_string(),
+ DataType::Map(MapType::new(
+ DataType::VarChar(VarCharType::string_type()),
+ DataType::Int(IntType::new()),
+ )),
+ ),
+ ];
+ let schema = build_target_arrow_schema(&fields).unwrap();
+
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ let output = file_io.new_output("memory:/reject_map.vortex").unwrap();
+
+ let err = match VortexFormatWriter::new(&output, schema).await {
+ Ok(_) => panic!("the vortex writer cannot convert a Map column"),
+ Err(err) => err,
+ };
+ let Error::Unsupported { message } = err else {
+ panic!("expected Unsupported, got {err:?}");
+ };
+ assert!(message.contains("'tags'"), "message: {message}");
+ assert!(message.contains("Map"), "message: {message}");
+ }
+
+ #[tokio::test]
+ async fn test_vortex_writer_rejects_a_multiset_nested_in_a_row() {
+ use crate::arrow::build_target_arrow_schema;
+ use crate::spec::{DataField, DataType, IntType, MultisetType, RowType,
VarCharType};
+
+ let fields = vec![DataField::new(
+ 0,
+ "nested".to_string(),
+ DataType::Row(RowType::new(vec![
+ DataField::new(1, "n".to_string(),
DataType::Int(IntType::new())),
+ DataField::new(
+ 2,
+ "m".to_string(),
+ DataType::Multiset(MultisetType::new(DataType::VarChar(
+ VarCharType::string_type(),
+ ))),
+ ),
+ ])),
+ )];
+ let schema = build_target_arrow_schema(&fields).unwrap();
+
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ let output = file_io
+ .new_output("memory:/reject_nested_multiset.vortex")
+ .unwrap();
+
+ let err = match VortexFormatWriter::new(&output, schema).await {
+ Ok(_) => panic!("a nested Multiset is just as unconvertible as a
top-level Map"),
+ Err(err) => err,
+ };
+ let Error::Unsupported { message } = err else {
+ panic!("expected Unsupported, got {err:?}");
+ };
+ assert!(message.contains("'nested'"), "message: {message}");
+ assert!(message.contains("Map"), "message: {message}");
+ }
+
+ /// Positive control: a nested ROW of primitives must still be accepted,
so the walk
+ /// cannot be satisfied by a blanket `Struct(..) => true`.
+ #[tokio::test]
+ async fn test_vortex_writer_accepts_a_nested_row_of_primitives() {
+ use crate::arrow::build_target_arrow_schema;
+ use crate::spec::{DataField, DataType, IntType, RowType, VarCharType};
+
+ let fields = vec![DataField::new(
+ 0,
+ "nested".to_string(),
+ DataType::Row(RowType::new(vec![
+ DataField::new(1, "n".to_string(),
DataType::Int(IntType::new())),
+ DataField::new(
+ 2,
+ "s".to_string(),
+ DataType::VarChar(VarCharType::string_type()),
+ ),
+ ])),
+ )];
+ let schema = build_target_arrow_schema(&fields).unwrap();
+
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ let output = file_io
+ .new_output("memory:/accept_nested_row.vortex")
+ .unwrap();
+
+ VortexFormatWriter::new(&output, schema)
+ .await
+ .expect("a nested ROW of primitives is convertible");
+ }
}