risingwave_connector/parser/
parquet_parser.rs1use std::sync::Arc;
16
17use futures_async_stream::try_stream;
18use prometheus::core::GenericCounter;
19use risingwave_common::array::arrow::arrow_array_iceberg::{ArrayRef, RecordBatch};
20use risingwave_common::array::arrow::arrow_schema_iceberg::FieldRef;
21use risingwave_common::array::arrow::{IcebergArrowConvert, is_parquet_field_match_source_schema};
22use risingwave_common::array::{ArrayBuilderImpl, DataChunk, StreamChunk};
23use risingwave_common::metrics::LabelGuardedMetric;
24use risingwave_common::types::{Datum, ScalarImpl};
25use thiserror_ext::AsReport;
26
27use crate::parser::ConnectorResult;
28use crate::source::SourceColumnDesc;
29#[derive(Debug)]
32pub struct ParquetParser {
33 rw_columns: Vec<SourceColumnDesc>,
34 file_name: String,
35 offset: usize,
36 case_insensitive: bool,
37}
38
39impl ParquetParser {
40 pub fn new(
41 rw_columns: Vec<SourceColumnDesc>,
42 file_name: String,
43 offset: usize,
44 case_insensitive: bool,
45 ) -> ConnectorResult<Self> {
46 Ok(Self {
47 rw_columns,
48 file_name,
49 offset,
50 case_insensitive,
51 })
52 }
53
54 #[try_stream(boxed, ok = StreamChunk, error = crate::error::ConnectorError)]
55 pub async fn into_stream(
56 mut self,
57 record_batch_stream: parquet::arrow::async_reader::ParquetRecordBatchStream<
58 tokio_util::compat::Compat<opendal::FuturesAsyncReader>,
59 >,
60 file_source_input_row_count_metrics: Option<
61 LabelGuardedMetric<GenericCounter<prometheus::core::AtomicU64>>,
62 >,
63 parquet_source_skip_row_count_metrics: Option<
64 LabelGuardedMetric<GenericCounter<prometheus::core::AtomicU64>>,
65 >,
66 ) {
67 #[for_await]
68 for record_batch in record_batch_stream {
69 let record_batch: RecordBatch = record_batch?;
70 let chunk: StreamChunk = self.convert_record_batch_to_stream_chunk(
72 record_batch,
73 file_source_input_row_count_metrics.clone(),
74 parquet_source_skip_row_count_metrics.clone(),
75 )?;
76
77 yield chunk;
78 }
79 }
80
81 fn inc_offset(&mut self) {
82 self.offset += 1;
83 }
84
85 fn convert_record_batch_to_stream_chunk(
106 &mut self,
107 record_batch: RecordBatch,
108 file_source_input_row_count_metrics: Option<
109 LabelGuardedMetric<GenericCounter<prometheus::core::AtomicU64>>,
110 >,
111 parquet_source_skip_row_count_metrics: Option<
112 LabelGuardedMetric<GenericCounter<prometheus::core::AtomicU64>>,
113 >,
114 ) -> Result<StreamChunk, crate::error::ConnectorError> {
115 const MAX_HIDDEN_COLUMN_NUMS: usize = 3;
116 let column_size = self.rw_columns.len();
117 let mut chunk_columns = Vec::with_capacity(self.rw_columns.len() + MAX_HIDDEN_COLUMN_NUMS);
118
119 for source_column in self.rw_columns.clone() {
120 match source_column.column_type {
121 crate::source::SourceColumnType::Normal => {
122 let rw_data_type: &risingwave_common::types::DataType =
123 &source_column.data_type;
124 let rw_column_name = &source_column.name;
125 if let Some((parquet_field, parquet_column)) =
126 self.find_parquet_column(&record_batch, rw_column_name)
127 && is_parquet_field_match_source_schema(parquet_field, rw_data_type)
128 {
129 let arrow_field = IcebergArrowConvert
133 .to_arrow_field(rw_column_name, rw_data_type)
134 .map_err(|e| {
135 crate::parser::AccessError::ParquetParser {
136 message: format!(
137 "to_arrow_field failed, column='{}', rw_type='{}', offset={}, error={}",
138 rw_column_name, rw_data_type, self.offset, e.as_report()
139 )
140 }
141 })?;
142 let array_impl = IcebergArrowConvert
143 .array_from_arrow_array(&arrow_field, parquet_column)
144 .map_err(|e| {
145 crate::parser::AccessError::ParquetParser {
146 message: format!(
147 "array_from_arrow_array failed, column='{}', rw_type='{}', arrow_field='{}', parquet_type='{}', offset={}, error={}",
148 rw_column_name,
149 rw_data_type,
150 arrow_field.data_type(),
151 parquet_column.data_type(),
152 self.offset,
153 e.as_report()
154 )
155 }
156 })?;
157 if array_impl.data_type() != *rw_data_type {
161 return Err(crate::parser::AccessError::ParquetParser {
162 message: format!(
163 "converted array type diverges from the declared column type, column='{}', rw_type='{}', converted='{}', parquet_type='{}', offset={}",
164 rw_column_name,
165 rw_data_type,
166 array_impl.data_type(),
167 parquet_column.data_type(),
168 self.offset,
169 ),
170 }
171 .into());
172 }
173 chunk_columns.push(Arc::new(array_impl));
174 } else {
175 let column = if let Some(additional_column_type) =
178 &source_column.additional_column.column_type
179 {
180 match additional_column_type {
181 risingwave_pb::plan_common::additional_column::ColumnType::Offset(_) => {
182 let mut array_builder = ArrayBuilderImpl::with_type(column_size, source_column.data_type.clone());
183 for _ in 0..record_batch.num_rows() {
184 let datum: Datum = Some(ScalarImpl::Utf8((self.offset).to_string().into()));
185 self.inc_offset();
186 array_builder.append(datum);
187 }
188 Arc::new(array_builder.finish())
189 }
190 risingwave_pb::plan_common::additional_column::ColumnType::Filename(_) => {
191 let mut array_builder = ArrayBuilderImpl::with_type(column_size, source_column.data_type.clone());
192 let datum: Datum = Some(ScalarImpl::Utf8(self.file_name.clone().into()));
193 array_builder.append_n(record_batch.num_rows(), datum);
194 Arc::new(array_builder.finish())
195 }
196 _ => unreachable!(),
197 }
198 } else {
199 let mut array_builder =
201 ArrayBuilderImpl::with_type(column_size, rw_data_type.clone());
202 array_builder.append_n_null(record_batch.num_rows());
203 if let Some(metrics) = parquet_source_skip_row_count_metrics.clone() {
204 metrics.inc_by(record_batch.num_rows() as u64);
205 }
206 Arc::new(array_builder.finish())
207 };
208 chunk_columns.push(column);
209 }
210 }
211 crate::source::SourceColumnType::RowId => {
212 let mut array_builder =
213 ArrayBuilderImpl::with_type(column_size, source_column.data_type.clone());
214 let datum: Datum = None;
215 array_builder.append_n(record_batch.num_rows(), datum);
216 let res = array_builder.finish();
217 let column = Arc::new(res);
218 chunk_columns.push(column);
219 }
220 crate::source::SourceColumnType::Offset | crate::source::SourceColumnType::Meta => {
222 unreachable!()
223 }
224 }
225 }
226 if let Some(metrics) = file_source_input_row_count_metrics {
227 metrics.inc_by(record_batch.num_rows() as u64);
228 }
229
230 let data_chunk = DataChunk::new(chunk_columns.clone(), record_batch.num_rows());
231 Ok(data_chunk.into())
232 }
233
234 fn find_parquet_column<'a>(
235 &self,
236 record_batch: &'a RecordBatch,
237 column_name: &str,
238 ) -> Option<(&'a FieldRef, &'a ArrayRef)> {
239 let fields = record_batch.schema_ref().fields();
240 if let Some(index) = fields.iter().position(|field| field.name() == column_name) {
241 return Some((&fields[index], record_batch.column(index)));
242 }
243 if !self.case_insensitive {
244 return None;
245 }
246 let mut matched_index: Option<usize> = None;
247 for (index, field) in fields.iter().enumerate() {
248 if field.name().eq_ignore_ascii_case(column_name) {
249 if matched_index.is_some() {
250 return None;
251 }
252 matched_index = Some(index);
253 }
254 }
255 matched_index.map(|index| (&fields[index], record_batch.column(index)))
256 }
257}