{"record":{"id":"c0009413baf733f1","repo":"risingwavelabs/risingwave","slug":"convert-parquet-batch-of-path","errorCode":null,"errorMessage":"convert parquet batch of {path}","messagePattern":"convert parquet batch of (.+?)","errorType":"exception","errorClass":"SinkError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/iceberg_with_pk_index/compaction_resolver.rs","lineNumber":434,"sourceCode":"    let metadata = input_file.metadata().await?;\n    let reader = input_file.reader().await?;\n    let builder = ParquetRecordBatchStreamBuilder::new(ParquetFileReader::new(metadata, reader))\n        .await\n        .map_err(|e| anyhow!(e).context(format!(\"open parquet reader for {path}\")))?;\n    let (projection, pk_order) = pk_projection(builder.parquet_schema(), pk_indices)?;\n    let mut stream = builder\n        .with_projection(projection)\n        .build()\n        .map_err(|e| anyhow!(e).context(format!(\"build parquet stream for {path}\")))?;\n\n    let mut base_pos = 0;\n    let mut iter = want_positions.iter().peekable();\n    while let Some(batch) = stream.next().await {\n        let batch =\n            batch.map_err(|e| anyhow!(e).context(format!(\"read parquet batch of {path}\")))?;\n        let chunk = IcebergArrowConvert\n            .chunk_from_record_batch(&batch)\n            .map_err(|e| anyhow!(e).context(format!(\"convert parquet batch of {path}\")))?\n            .project(&pk_order);\n        while let Some(&iter_pos) = iter.peek() {\n            let chunk_pos = iter_pos as usize - base_pos;\n            if chunk_pos >= chunk.capacity() {\n                break;\n            }\n            let pk = chunk.row_at(chunk_pos).0.to_owned_row();\n            results.push(pk);\n            iter.next();\n        }\n        base_pos += chunk.capacity();\n    }\n\n    Ok(results)\n}\n\n#[try_stream(ok = DataChunk, error = SinkError)]\nasync fn scan_output_file_inner<'a>(","sourceCodeStart":416,"sourceCodeEnd":452,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/iceberg_with_pk_index/compaction_resolver.rs#L416-L452","documentation":"Thrown in `scan_input_pks_at_positions` when `IcebergArrowConvert::chunk_from_record_batch` cannot convert a decoded Arrow `RecordBatch` into a RisingWave `DataChunk`. The parquet read succeeded, but the Arrow batch does not match what the converter expects (unsupported data type in a column, mismatched schema vs. expectation, or nullability/layout issues), so the batch cannot be projected onto the pk columns.","triggerScenarios":"`resolve` → `scan_input_pks_at_positions` converts each parquet batch to a chunk and projects it with `pk_order`. Triggered when the parquet file contains Arrow types the IcebergArrowConvert does not support (e.g. exotic logical types, dictionary-encoded columns the converter rejects), or the projected batch shape does not match the expected schema.","commonSituations":"Files written by an external Iceberg writer using data types RisingWave's converter does not yet support; schema evolution introduced a new type into pk-adjacent columns; parquet logical/converted-type written inconsistently by a different writer version.","solutions":["Inspect the conversion error to find the offending column/type; check RisingWave's supported Iceberg type mappings.","Rewrite the offending file with a writer that emits supported types (e.g. via Iceberg rewrite action).","Verify pk columns use plain supported types (int/long/string/etc.) in the table schema.","If a RisingWave version issue, upgrade to a version with broader IcebergArrowConvert type support."],"exampleFix":"// before: table schema with unsupported type in a column reached during projection\n-- column updated_at TIMESTAMPNSTZ (unsupported)\n\n// after\n-- column updated_at TIMESTAMPTZ or TIMESTAMP (supported mapping)","handlingStrategy":"try-catch","validationCode":"fn supports_arrow_types(schema: &arrow_schema::SchemaRef) -> bool {\n    schema.fields().iter().all(|f| matches!(f.data_type(),\n        arrow_schema::DataType::Boolean\n        | arrow_schema::DataType::Int32 | arrow_schema::DataType::Int64\n        | arrow_schema::DataType::Float32 | arrow_schema::DataType::Float64\n        | arrow_schema::DataType::Utf8 | arrow_schema::DataType::Binary\n        | arrow_schema::DataType::Timestamp(_, _)))\n}","typeGuard":null,"tryCatchPattern":"let chunk = match IcebergArrowConvert.chunk_from_record_batch(&batch) {\n    Ok(c) => c,\n    Err(e) => return Err(SinkError::Iceberg(anyhow!(e).context(\"unsupported arrow type in parquet batch\"))),\n};","preventionTips":["Restrict sink table schemas to types the IcebergArrowConvert supports.","Test round-trip conversion (chunk -> parquet -> chunk) when adding new column types.","Keep writer and reader RisingWave versions aligned during rolling upgrades."],"tags":["parquet","arrow","type-conversion","iceberg-sink"],"backgroundTag":"type-mismatch","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}