Package org.apache.iceberg.flink.data
Class FlinkVariantShreddingAnalyzer
java.lang.Object
org.apache.iceberg.parquet.VariantShreddingAnalyzer<org.apache.flink.table.data.RowData,org.apache.flink.table.types.logical.RowType>
org.apache.iceberg.flink.data.FlinkVariantShreddingAnalyzer
public class FlinkVariantShreddingAnalyzer
extends VariantShreddingAnalyzer<org.apache.flink.table.data.RowData,org.apache.flink.table.types.logical.RowType>
Analyzes Variant fields in Flink
RowData and converts Flink's binary Variant
representation to Iceberg VariantValue instances for Variant shredding.-
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionprotected List<VariantValue>extractVariantValues(List<org.apache.flink.table.data.RowData> bufferedRows, int variantFieldIndex) protected intresolveColumnIndex(org.apache.flink.table.types.logical.RowType flinkSchema, String columnName) Resolves a column name to its index in the engine-specific schema.Methods inherited from class org.apache.iceberg.parquet.VariantShreddingAnalyzer
analyzeAndCreateSchema, analyzeVariantColumns
-
Constructor Details
-
FlinkVariantShreddingAnalyzer
public FlinkVariantShreddingAnalyzer()
-
-
Method Details
-
extractVariantValues
protected List<VariantValue> extractVariantValues(List<org.apache.flink.table.data.RowData> bufferedRows, int variantFieldIndex) - Specified by:
extractVariantValuesin classVariantShreddingAnalyzer<org.apache.flink.table.data.RowData,org.apache.flink.table.types.logical.RowType>
-
resolveColumnIndex
protected int resolveColumnIndex(org.apache.flink.table.types.logical.RowType flinkSchema, String columnName) Description copied from class:VariantShreddingAnalyzerResolves a column name to its index in the engine-specific schema. Returns -1 if the column is not found.- Specified by:
resolveColumnIndexin classVariantShreddingAnalyzer<org.apache.flink.table.data.RowData,org.apache.flink.table.types.logical.RowType>
-