Class HudiEngine
java.lang.Object
com.logicalclocks.hsfs.spark.engine.hudi.HudiEngine
-
Field Summary
FieldsModifier and TypeFieldDescriptionprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringstatic final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static HudiEngineprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final Stringprotected static final String -
Method Summary
Modifier and TypeMethodDescriptiondeleteRecord(org.apache.spark.sql.SparkSession sparkSession, FeatureGroupBase featureGroup, org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> deleteDF, Map<String, String> writeOptions) getHudiTableSize(org.apache.spark.api.java.JavaSparkContext jsc, String basePath) static HudiEnginevoidreconcileHudiSchema(org.apache.spark.sql.SparkSession sparkSession, FeatureGroupAlias featureGroupAlias, Map<String, String> hudiArgs) voidsaveHudiFeatureGroup(org.apache.spark.sql.SparkSession sparkSession, FeatureGroupBase featureGroup, org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> dataset, HudiOperationType operation, Map<String, String> writeOptions, Integer validationId) booleansparkSchemasMatch(String[] schema1, String[] schema2) voidstreamToHoodieTable(org.apache.spark.sql.SparkSession sparkSession, StreamFeatureGroup streamFeatureGroup, Map<String, String> writeOptions)
-
Field Details
-
HUDI_SPARK_FORMAT
- See Also:
-
HUDI_TABLE_VERSION
- See Also:
-
HUDI_BASE_PATH
- See Also:
-
HUDI_TABLE_NAME
- See Also:
-
HUDI_TABLE_TYPE
- See Also:
-
HUDI_TABLE_STORAGE_TYPE
- See Also:
-
HUDI_TABLE_OPERATION
- See Also:
-
HUDI_METADATA_ENABLE
- See Also:
-
HUDI_TABLE_RECORD_KEY_FIELD
- See Also:
-
HUDI_TABLE_PARTITION_KEY_FIELDS
- See Also:
-
HUDI_TABLE_KEY_GENERATOR_CLASS
- See Also:
-
HUDI_TABLE_PRECOMBINE_FIELD
- See Also:
-
HUDI_TABLE_BASE_FILE_FORMAT
- See Also:
-
HUDI_TABLE_METADATA_PARTITIONS
- See Also:
-
HUDI_KEY_GENERATOR_OPT_KEY
- See Also:
-
HUDI_COMPLEX_KEY_GENERATOR_OPT_VAL
- See Also:
-
HUDI_WRITE_RECORD_KEY
- See Also:
-
HUDI_PARTITION_FIELD
- See Also:
-
HUDI_WRITE_PRECOMBINE_FIELD
- See Also:
-
HUDI_HIVE_SYNC_ENABLE
- See Also:
-
HUDI_HIVE_SYNC_TABLE
- See Also:
-
HUDI_HIVE_SYNC_DB
- See Also:
-
HUDI_HIVE_SYNC_MODE
- See Also:
-
HUDI_HIVE_SYNC_MODE_VAL
- See Also:
-
HUDI_HIVE_SYNC_PARTITION_FIELDS
- See Also:
-
HIVE_PARTITION_EXTRACTOR_CLASS_OPT_KEY
- See Also:
-
DEFAULT_HIVE_PARTITION_EXTRACTOR_CLASS_OPT_VAL
- See Also:
-
HIVE_NON_PARTITION_EXTRACTOR_CLASS_OPT_VAL
- See Also:
-
HIVE_AUTO_CREATE_DATABASE_OPT_KEY
- See Also:
-
HIVE_AUTO_CREATE_DATABASE_OPT_VAL
- See Also:
-
HUDI_COPY_ON_WRITE
- See Also:
-
HUDI_QUERY_TYPE_OPT_KEY
- See Also:
-
HUDI_QUERY_TYPE_INCREMENTAL_OPT_VAL
- See Also:
-
HUDI_QUERY_TYPE_SNAPSHOT_OPT_VAL
- See Also:
-
HUDI_QUERY_TIME_TRAVEL_AS_OF_INSTANT
- See Also:
-
HUDI_BEGIN_INSTANTTIME_OPT_KEY
- See Also:
-
HUDI_END_INSTANTTIME_OPT_KEY
- See Also:
-
HUDI_WRITE_INSERT_DROP_DUPLICATES
- See Also:
-
HUDI_KAFKA_TOPIC
- See Also:
-
COMMIT_METADATA_KEYPREFIX_OPT_KEY
- See Also:
-
STREAMER_CHECKPOINT_KEY_V2
- See Also:
-
DELTASTREAMER_CHECKPOINT_KEY
- See Also:
-
INITIAL_CHECKPOINT_STRING
- See Also:
-
FEATURE_GROUP_SCHEMA
- See Also:
-
FEATURE_GROUP_ENCODED_SCHEMA
- See Also:
-
FEATURE_GROUP_COMPLEX_FEATURES
- See Also:
-
KAFKA_SOURCE
- See Also:
-
SCHEMA_PROVIDER
- See Also:
-
DELTA_STREAMER_TRANSFORMER
- See Also:
-
DELTA_SOURCE_ORDERING_FIELD_OPT_KEY
- See Also:
-
MIN_SYNC_INTERVAL_SECONDS
- See Also:
-
SPARK_MASTER
- See Also:
-
PROJECT_ID
- See Also:
-
FEATURE_STORE_NAME
- See Also:
-
SUBJECT_ID
- See Also:
-
FEATURE_GROUP_ID
- See Also:
-
FEATURE_GROUP_NAME
- See Also:
-
FEATURE_GROUP_VERSION
- See Also:
-
FUNCTION_TYPE
- See Also:
-
STREAMING_QUERY
- See Also:
-
INSTANCE
-
-
Method Details
-
getInstance
-
saveHudiFeatureGroup
public void saveHudiFeatureGroup(org.apache.spark.sql.SparkSession sparkSession, FeatureGroupBase featureGroup, org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> dataset, HudiOperationType operation, Map<String, String> writeOptions, Integer validationId) throws IOException, FeatureStoreException, ParseException -
deleteRecord
public FeatureGroupCommit deleteRecord(org.apache.spark.sql.SparkSession sparkSession, FeatureGroupBase featureGroup, org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> deleteDF, Map<String, String> writeOptions) throws IOException, FeatureStoreException, ParseException -
setupHudiReadOpts
-
reconcileHudiSchema
public void reconcileHudiSchema(org.apache.spark.sql.SparkSession sparkSession, FeatureGroupAlias featureGroupAlias, Map<String, String> hudiArgs) throws FeatureStoreException- Throws:
FeatureStoreException
-
sparkSchemasMatch
-
streamToHoodieTable
public void streamToHoodieTable(org.apache.spark.sql.SparkSession sparkSession, StreamFeatureGroup streamFeatureGroup, Map<String, String> writeOptions) throws Exception- Throws:
Exception
-
getHudiTableSize
public Long getHudiTableSize(org.apache.spark.api.java.JavaSparkContext jsc, String basePath) throws IOException - Throws:
IOException
-