Writing Data to Hudi Tables In Batches
Scenarios
Hudi provides multiple write modes. For details, see the configuration item hoodie.datasource.write.operation. This section describes upsert, insert, and bulk_insert.
- insert: The operation process is similar to upsert. The query on updated file partitions is not based on indexes. Therefore, insert is faster than upsert. This operation is recommended for data sources that do not contain updated data. If the data source contains updated data, the data lake will have duplicate data.
- bulk_insert (insert in batches): It is used for initial dataset loading. This operation sorts primary keys and then inserts data into a Hudi table by writing data to a common Parquet table. It has the best performance but cannot control small files. The upsert and insert operations can control small files by using heuristics.
- upsert (insert and update): It is the default operation type. Hudi determines whether historical data exists based on the primary key. Historical data is updated, and other data is inserted. This operation is recommended for data sources, such as change data capture (CDC), that include updated data.
- Primary keys are not sorted during insert. Therefore, you are not advised to use insert during dataset initialization.
- You are advised to use insert if data is new, use upsert if data needs to be updated, and use bulk_insert if datasets need to be initialized.
- If the inserted or upserted data exceeds the defined DECIMAL type, the system reports an error, indicating data too long. While performing this data type validation may slightly reduce write speeds, it ensures data validity by reporting an error if the input does not match the defined data type.
- The bulk_insert mode bypasses data type validation to maximize throughput. This allows values exceeding the DECIMAL precision or scale to be written, potentially leading to data corruption. You must ensure all batch inputs are valid before being written.
- Performing a bulk_insert on a large number of partitions simultaneously may lead to OOM errors. To prevent this, enable global sorting to optimize memory usage during the write process. For SQL, use set hoodie.bulkinsert.sort.mode = GLOBAL_SORT. For APIs, use option("hoodie.bulkinsert.sort.mode", "GLOBAL_SORT").
Writing Data to Hudi Tables In Batches
- Import the Hudi package to generate test data. For details, see 1 to 3 in Creating a Hudi Table Using Spark Shell.
- Write data to the Hudi table. Add the option("hoodie.datasource.write.operation", "bulk_insert") parameter to the write command and set the write mode to bulk_insert. For details about how to specify other write modes, see Table 1.
- Currently, bulk insert of Hudi is implemented using Spark V2 APIs. If a Hive-incompatible bulk insert format is used to write data, an error may be reported. For details, see Bulk Insert Failed to Write Data, with the Error Message "LEGACY store assignment policy is disallowed in Spark data source V2" Displayed.
- By default, bulk insert cannot be used for bucket indexes. To use bulk insert, configure option("hoodie.bucket.support.bulk.insert", true) and option("hoodie.datasource.write.row.writer.enable", true). A table with bucket indexes can be written only once using bulk insert.
df.write.format("org.apache.hudi"). options(getQuickstartWriteConfigs). option("hoodie.datasource.write.precombine.field", "ts"). option("hoodie.datasource.write.recordkey.field", "uuid"). option("hoodie.datasource.write.partitionpath.field", ""). option("hoodie.datasource.write.operation", "bulk_insert"). option("hoodie.table.name", tableName). option("hoodie.datasource.write.keygenerator.class", "org.apache.hudi.keygen.NonpartitionedKeyGenerator"). option("hoodie.datasource.hive_sync.enable", "true"). option("hoodie.datasource.hive_sync.partition_fields", ""). option("hoodie.datasource.hive_sync.partition_extractor_class", "org.apache.hudi.hive.NonPartitionedExtractor"). option("hoodie.datasource.hive_sync.table", tableName). option("hoodie.datasource.hive_sync.use_jdbc", "false"). option("hoodie.bulkinsert.shuffle.parallelism", 4). mode(Overwrite). save(basePath)
- For details about the parameters in the example, see Table 1.
- If the Spark DataSource API is used to update the MOR table, small files of the updated data may be merged when a small volume of data is inserted. As a result, some updated data can be found in the read-optimized view of the MOR table.
- If the base file of the data to be updated is a small file, the data to be inserted and new data for update are merged with the base file to generate a new base file instead of being written to logs.
Configuring Partitions
Hudi supports multiple partitioning modes, such as multi-level partitioning, non-partitioning, single-level partitioning, and partitioning by date. You can select a proper partitioning mode as required. The following describes how to configure different partitioning modes for Hudi.
- Multi-level partitioning
Multi-level partitioning indicates that multiple fields are specified as partition keys. Pay attention to the following configuration items:
Configuration Item
Description
hoodie.datasource.write.partitionpath.field
Data partition field. Set this parameter to partition data by specific fields, improving query performance and data management.
Configure multiple partition fields, for example, p1, p2, and p3.
hoodie.datasource.hive_sync.partition_fields
Partition field for synchronizing Hudi tables with the Hive metadata.
Set this parameter to the same value as hoodie.datasource.write.partitionpath.field.
hoodie.datasource.write.keygenerator.class
Key generator class used to generate record keys and partition paths in Hudi tables. This generator ensures data uniqueness and enables efficient partition management.
Set this parameter to org.apache.hudi.keygen.ComplexKeyGenerator.
hoodie.datasource.hive_sync.partition_extractor_class
Class used to specify the partition extractor. The partition extractor extracts partition values from the Hudi table's partition path. This ensures that partitions are correctly created when Hudi tables are synchronized to Hive.
Set this parameter to org.apache.hudi.hive.MultiPartKeysValueExtractor.
df.write.format("org.apache.hudi"). options(getQuickstartWriteConfigs). option("hoodie.datasource.write.precombine.field", "ts"). option("hoodie.datasource.write.recordkey.field", "uuid"). option("hoodie.datasource.write.partitionpath.field", "p1,p2,p3"). option("hoodie.datasource.write.operation", "bulk_insert"). option("hoodie.table.name", tableName). option("hoodie.datasource.write.keygenerator.class", "org.apache.hudi.keygen.ComplexKeyGenerator"). option("hoodie.datasource.hive_sync.enable", "true"). option("hoodie.datasource.hive_sync.partition_fields", "p1,p2,p3"). option("hoodie.datasource.hive_sync.partition_extractor_class", "org.apache.hudi.hive.MultiPartKeysValueExtractor"). option("hoodie.datasource.hive_sync.table", tableName). option("hoodie.datasource.hive_sync.use_jdbc", "false"). option("hoodie.bulkinsert.shuffle.parallelism", 4). mode(Overwrite). save(basePath)
- Non-partitioning
Hudi supports non-partitioned tables. Pay attention to the following configuration items:
Configuration Item
Description
hoodie.datasource.write.partitionpath.field
Leave this parameter blank.
hoodie.datasource.hive_sync.partition_fields
Leave this parameter blank.
hoodie.datasource.write.keygenerator.class
Set this parameter to org.apache.hudi.keygen.NonpartitionedKeyGenerator.
hoodie.datasource.hive_sync.partition_extractor_class
Set this parameter to org.apache.hudi.hive.NonPartitionedExtractor.
df.write.format("org.apache.hudi"). options(getQuickstartWriteConfigs). option("hoodie.datasource.write.precombine.field", "ts"). option("hoodie.datasource.write.recordkey.field", "uuid"). option("hoodie.datasource.write.partitionpath.field", ""). option("hoodie.datasource.write.operation", "bulk_insert"). option("hoodie.table.name", tableName). option("hoodie.datasource.write.keygenerator.class", "org.apache.hudi.keygen.NonpartitionedKeyGenerator"). option("hoodie.datasource.hive_sync.enable", "true"). option("hoodie.datasource.hive_sync.partition_fields", ""). option("hoodie.datasource.hive_sync.partition_extractor_class", "org.apache.hudi.hive.NonPartitionedExtractor"). option("hoodie.datasource.hive_sync.table", tableName). option("hoodie.datasource.hive_sync.use_jdbc", "false"). option("hoodie.bulkinsert.shuffle.parallelism", 4). mode(Overwrite). save(basePath)
- Single-level partitioning
It is similar to multi-level partitioning. Pay attention to the following configuration items:
Configuration Item
Description
hoodie.datasource.write.partitionpath.field
Set this parameter to one field, for example, p.
hoodie.datasource.hive_sync.partition_fields
Set this same as hoodie.datasource.write.partitionpath.field.
hoodie.datasource.write.keygenerator.class
By default, you can set this parameter to org.apache.hudi.keygen.SimpleKeyGenerator and org.apache.hudi.keygen.ComplexKeyGenerator. You can also leave this parameter blank.
hoodie.datasource.hive_sync.partition_extractor_class
Set this parameter to org.apache.hudi.hive.MultiPartKeysValueExtractor.
df.write.format("org.apache.hudi"). options(getQuickstartWriteConfigs). option("hoodie.datasource.write.precombine.field", "ts"). option("hoodie.datasource.write.recordkey.field", "uuid"). option("hoodie.datasource.write.partitionpath.field", "p"). option("hoodie.datasource.write.operation", "bulk_insert"). option("hoodie.table.name", tableName). option("hoodie.datasource.write.keygenerator.class", "org.apache.hudi.keygen.ComplexKeyGenerator"). option("hoodie.datasource.hive_sync.enable", "true"). option("hoodie.datasource.hive_sync.partition_fields", "p"). option("hoodie.datasource.hive_sync.partition_extractor_class", "org.apache.hudi.hive.MultiPartKeysValueExtractor"). option("hoodie.datasource.hive_sync.table", tableName). option("hoodie.datasource.hive_sync.use_jdbc", "false"). option("hoodie.bulkinsert.shuffle.parallelism", 4). mode(Overwrite). save(basePath)
- Partitioning by date
The date field is specified as the partition field. Pay attention to the following configuration items:
Configuration Item
Description
hoodie.datasource.write.partitionpath.field
Set this parameter to the date field.
hoodie.datasource.hive_sync.partition_fields
Set this same as hoodie.datasource.write.partitionpath.field.
hoodie.datasource.write.keygenerator.class
Set this parameter to org.apache.hudi.keygen.ComplexKeyGenerator.
hoodie.datasource.hive_sync.partition_extractor_class
Set this parameter to org.apache.hudi.hive.SlashEncodedDayPartitionValueExtractor.
Date format for SlashEncodedDayPartitionValueExtractor must be yyyy/mm/dd.
- Partition sorting
Configuration Item
Description
hoodie.bulkinsert.user.defined.partitioner.class
Class used to specify the partition extractor. The partition extractor extracts partition values from the Hudi table's partition path. This ensures that partitions are correctly created when Hudi tables are synchronized to Hive.
Hudi provides multiple partition extractor classes, each supporting a different partition path format.
- MultiPartKeysValueExtractor: The default partition extractor. It is used for multi-level partition paths.
- Class name: org.apache.hudi.hive.MultiPartKeysValueExtractor
- Applicable scenario: The partition path format is part1=value1/part2=value2/....
- NonPartitionedExtractor: Used for tables that do not contain partitions.
- Class name: org.apache.hudi.hive.NonPartitionedExtractor
- Applicable scenario: The table does not have a partition path.
- HiveStylePartitionValueExtractor: Used for Hive-style partition paths.
- Class name: org.apache.hudi.hive.HiveStylePartitionValueExtractor
- Applicable scenario: The partition path format is part1=value1/part2=value2/..., which matches the Hive partition path format.
- CustomPartitionValueExtractor: A user-defined partition extractor class. This is used for custom requirements.
- Class name: User-defined class name
- Applicable scenario: Used to extract partition values according to custom requirements.
By default, bulk_insert sorts data by character and applies only to primary keys of StringType.
- MultiPartKeysValueExtractor: The default partition extractor. It is used for multi-level partition paths.
What is your overall rating for this page?
Thank you very much for your feedback. We will continue working to improve the documentation.See the reply and handling status in My Cloud VOC.
For any further questions, feel free to contact us through the chatbot.
Chatbot