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 Details

    • FlinkVariantShreddingAnalyzer

      public FlinkVariantShreddingAnalyzer()
  • Method Details

    • extractVariantValues

      protected List<VariantValue> extractVariantValues(List<org.apache.flink.table.data.RowData> bufferedRows, int variantFieldIndex)
      Specified by:
      extractVariantValues in class VariantShreddingAnalyzer<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: VariantShreddingAnalyzer
      Resolves a column name to its index in the engine-specific schema. Returns -1 if the column is not found.
      Specified by:
      resolveColumnIndex in class VariantShreddingAnalyzer<org.apache.flink.table.data.RowData,org.apache.flink.table.types.logical.RowType>