Class DynamicTableRecordGenerator

java.lang.Object
org.apache.iceberg.flink.sink.dynamic.DynamicTableRecordGenerator
All Implemented Interfaces:
Serializable, DynamicRecordGenerator<org.apache.flink.table.data.RowData>
Direct Known Subclasses:
VariantAvroDynamicTableRecordGenerator

public abstract class DynamicTableRecordGenerator extends Object implements DynamicRecordGenerator<org.apache.flink.table.data.RowData>
Abstract base class for SQL-based dynamic record generators. Users will extend this class to create a DynamicRecord from RowData.
See Also:
  • Constructor Details

    • DynamicTableRecordGenerator

      public DynamicTableRecordGenerator(org.apache.flink.table.types.logical.RowType rowType)
    • DynamicTableRecordGenerator

      public DynamicTableRecordGenerator(org.apache.flink.table.types.logical.RowType rowType, Map<String,String> writeProperties, org.apache.flink.configuration.Configuration flinkConfiguration)
  • Method Details

    • open

      public void open(org.apache.flink.api.common.functions.OpenContext openContext) throws Exception
      Specified by:
      open in interface DynamicRecordGenerator<org.apache.flink.table.data.RowData>
      Throws:
      Exception
    • rowType

      protected org.apache.flink.table.types.logical.RowType rowType()
    • flinkDynamicSinkConf

      protected FlinkDynamicSinkConf flinkDynamicSinkConf()
    • fieldNameToPosition

      protected Map<String,Integer> fieldNameToPosition()
    • validateRequiredColumnAndType

      protected void validateRequiredColumnAndType(String columnName, org.apache.flink.table.types.logical.LogicalType expectedType)