Insert or overwrite data into dynamic partitions

Updated at:

MaxCompute lets you insert data into dynamic partitions with the INSERT INTO or INSERT OVERWRITE statement.

Prerequisites

To perform INSERT INTO and INSERT OVERWRITE operations, you must have Update permission on the target table and Select permission on the source table. For more information about how to grant permissions, see MaxCompute Permissions.

How it works

When you use MaxCompute SQL to process data, you do not directly specify partition values in an INSERT INTO or INSERT OVERWRITE statement. You specify only the partition columns. The values for the partition columns are provided in the SELECT clause, and the system automatically inserts the data into the corresponding partitions based on these values.

For information about inserting data into static partitions, see Insert or overwrite data (INSERT INTO | INSERT OVERWRITE).

Limits

The limitations on inserting data into dynamic partitions by using the INSERT INTO and INSERT OVERWRITE operations are as follows:

  • An INSERT INTO statement can generate up to 10,000 dynamic partitions. An INSERT OVERWRITE statement can generate up to 60,000 dynamic partitions.

  • In a distributed environment, a single process that runs an SQL statement for a dynamic partition operation can write to up to 512 dynamic partitions. Otherwise, the job fails.

  • A dynamically generated partition value cannot be NULL. It also cannot contain special characters or Chinese characters. Otherwise, the following error is reported: FAILED: ODPS-0123031:Partition exception - invalid dynamic partition value: province=xxx.

    Note

    The value of a partition key column cannot contain double-byte characters, such as Chinese characters. The value of a partition key column must start with a letter and can contain letters, digits, and supported special characters. It must be 1 to 255 bytes in length. The following special characters are supported: spaces, colons (:), underscores (_), dollar signs ($), number signs (#), periods (.), exclamation points (!), and at signs (@). The behavior of other characters is not defined, such as escape characters \t, \n, and /.

  • Clustered tables do not support dynamic partitions.

Usage notes

When you insert data into a dynamic partition, note the following:

  • If the destination partition does not exist when you run an INSERT INTO PARTITION statement, MaxCompute automatically creates it.

  • If multiple concurrent INSERT INTO PARTITION jobs attempt to write to a non-existent partition, the first successful job creates the partition. Only one partition is created.

  • If you cannot control the concurrency of INSERT INTO PARTITION jobs, create partitions in advance with the ALTER TABLE command. For more information, see Partition operations.

  • If the destination table has multiple partition levels, you can specify some partitions as static partitions in an INSERT statement. However, static partitions must be higher-level partitions than any dynamic partitions.

  • When you insert data into a dynamic partition, the dynamic partition columns must be included in the SELECT list. Otherwise, the statement fails.

Syntax

INSERT {INTO|OVERWRITE} TABLE <table_name> PARTITION (<ptcol_name>[, <ptcol_name> ...]) 
<select_statement> FROM <from_statement>;

Parameters

Parameter

Required

Description

table_name

Yes

The name of the destination table.

ptcol_name

Yes

The name of a partition column in the destination table.

select_statement

Yes

The SELECT clause that queries data from the source table to be inserted into the destination table.

The mapping between the columns in the SELECT list and the dynamic partitions is determined by column order, not by name. The last columns in the SELECT list provide the values for the dynamic partition columns. For example, if a table has one dynamic partition, the last column in the SELECT list is used as the partition value. If the column order of the source table differs from that of the destination table, explicitly list the columns in your select_statement to ensure the correct order.

from_statement

Yes

The FROM clause that specifies the data source, such as the name of the source table.

Sample data

-- Create a partitioned table named sale_detail.
CREATE TABLE IF NOT EXISTS sale_detail
(
shop_name     STRING,
customer_id   STRING,
total_price   DOUBLE
)
PARTITIONED BY (sale_date STRING, region STRING);

-- Add a partition to the source table. This step is optional. If you do not create the partition in advance, it is automatically created during the data write operation.
ALTER TABLE sale_detail ADD PARTITION (sale_date='2013', region='china');

-- Append data to the source table. You can omit the TABLE keyword after INSERT INTO and INSERT OVERWRITE.
INSERT INTO sale_detail PARTITION (sale_date='2013', region='china') VALUES ('s1','c1',100.1),('s2','c2',100.2),('s3','c3',100.3);

-- Enable a full scan. This setting is effective only for the current session. Run a SELECT statement to view the data in the sale_detail table.
SET odps.sql.allow.fullscan=true; 
SELECT * FROM sale_detail;

-- Output:
+------------+-------------+-------------+------------+------------+
| shop_name  | customer_id | total_price | sale_date  | region     |
+------------+-------------+-------------+------------+------------+
| s1         | c1          | 100.1       | 2013       | china      |
| s2         | c2          | 100.2       | 2013       | china      |
| s3         | c3          | 100.3       | 2013       | china      |
+------------+-------------+-------------+------------+------------+

Examples

The following examples use the data from the sale_detail table created in the Sample data section.

  • This example inserts data into a destination table. The resulting dynamic partitions are determined by the values in the region column.

    -- Create a destination table named total_revenues.
    CREATE TABLE total_revenues (revenue DOUBLE) PARTITIONED BY (region string);
    
    -- Insert data from the sale_detail source table into the total_revenues destination table.
    SET odps.sql.allow.fullscan=true; 
    INSERT OVERWRITE TABLE total_revenues PARTITION(region) SELECT total_price AS revenue, region FROM sale_detail;
    
    -- Run the SHOW PARTITIONS statement to view the partitions of the total_revenues table.
    SHOW PARTITIONS total_revenues;
       
    -- Output:
    region=china  
    
    -- Enable a full scan. This setting is effective only for the current session. Run a SELECT statement to view the data in the total_revenues table. 
    SET odps.sql.allow.fullscan=true; 
    SELECT * FROM total_revenues;    
    
    -- Output:
    +------------+------------+
    | revenue    | region     |
    +------------+------------+
    | 100.1      | china      |
    | 100.2      | china      |
    | 100.3      | china      |
    +------------+------------+        
  • This example inserts data into a multi-level partitioned table and specifies the value for the first-level partition, sale_date.

    -- Create a destination table named sale_detail_dypart. 
    CREATE TABLE sale_detail_dypart LIKE sale_detail; 
    
    -- Specify the first-level partition and insert data into the destination table.
    SET odps.sql.allow.fullscan=true; 
    INSERT OVERWRITE TABLE sale_detail_dypart PARTITION (sale_date='2013', region)
    SELECT shop_name,customer_id,total_price,region FROM sale_detail;
    
    -- Enable a full scan. This setting is effective only for the current session. Run a SELECT statement to view the data in the sale_detail_dypart table.
    SET odps.sql.allow.fullscan=true; 
    SELECT * FROM sale_detail_dypart;
    
    -- Output:
    +------------+-------------+-------------+------------+------------+
    | shop_name  | customer_id | total_price | sale_date  | region     |
    +------------+-------------+-------------+------------+------------+
    | s1         | c1          | 100.1       | 2013       | china      |
    | s2         | c2          | 100.2       | 2013       | china      |
    | s3         | c3          | 100.3       | 2013       | china      |
    +------------+-------------+-------------+------------+------------+
  • Example 3: In dynamic partitioning, the mapping between the fields of the select_statement and the dynamic partitions of the destination table is determined by the order of the fields, not by the column names. The following is a sample command:

    -- Insert data from the source table sale_detail into the destination table sale_detail_dypart.
    SET odps.sql.allow.fullscan=true; 
    INSERT OVERWRITE TABLE sale_detail_dypart PARTITION (sale_date, region)
    SELECT shop_name,customer_id,total_price,sale_date,region FROM sale_detail;
    
    -- Enable a full scan. This setting is effective only for the current session. Run a SELECT statement to view the data in the sale_detail_dypart table.
    SET odps.sql.allow.fullscan=true; 
    SELECT * FROM sale_detail_dypart;
    
    -- Output:
    -- The sale_date field from the source table maps to the sale_date dynamic partition column of the destination table.
    -- The region field from the source table maps to the region dynamic partition column of the destination table.
    +------------+-------------+-------------+------------+------------+
    | shop_name  | customer_id | total_price | sale_date  | region     |
    +------------+-------------+-------------+------------+------------+
    | s1         | c1          | 100.1       | 2013       | china      |
    | s2         | c2          | 100.2       | 2013       | china      |
    | s3         | c3          | 100.3       | 2013       | china      |
    +------------+-------------+-------------+------------+------------+
    
    
    -- Insert data from the sale_detail table into the sale_detail_dypart table and change the SELECT field order.
    SET odps.sql.allow.fullscan=true; 
    INSERT OVERWRITE TABLE sale_detail_dypart PARTITION (sale_date, region)
    SELECT shop_name,customer_id,total_price,region,sale_date FROM sale_detail;
    
    -- Enable a full scan. This setting is effective only for the current session. Run a SELECT statement to view the data in the sale_detail_dypart table.
    SET odps.sql.allow.fullscan=true; 
    SELECT * FROM sale_detail_dypart;
    
    -- Output:
    -- The region field from the source table maps to the sale_date dynamic partition column of the destination table.
    -- The sale_date field from the source table maps to the region dynamic partition column of the destination table.
    +------------+-------------+-------------+------------+------------+
    | shop_name  | customer_id | total_price | sale_date  | region     |
    +------------+-------------+-------------+------------+------------+
    | s1         | c1          | 100.1       | china      | 2013       |
    | s2         | c2          | 100.2       | china      | 2013       |
    | s3         | c3          | 100.3       | china      | 2013       |
    +------------+-------------+-------------+------------+------------+
  • Example 4 (Incorrect): This example fails because the dynamic partition column is not included in the SELECT list.

    INSERT OVERWRITE TABLE sale_detail_dypart PARTITION (sale_date='2013', region) 
    SELECT shop_name,customer_id,total_price FROM sale_detail;

    Output:

    FAILED: ODPS-0130071:[1,24] Semantic analysis exception - wrong columns count 3 in data source, requires 4 columns (includes dynamic partitions if any)
  • Example 5 (Incorrect): This example fails because a higher-level partition cannot be dynamic if a lower-level partition is static.

    INSERT OVERWRITE TABLE sale_detail_dypart PARTITION (sale_date, region='china')
    SELECT shop_name,customer_id,total_price,sale_date FROM sale_detail_dypart;

    Output:

    FAILED: ODPS-0130071:[1,72] Semantic analysis exception - static partition region must be a high level partition than any dynamic partitions
  • Example 6: When inserting data into a dynamic partition, MaxCompute performs an implicit conversion if the data type of the partition column does not strictly match that of the corresponding column in the SELECT statement. The following is an example command:

    -- Create a source table named src.
    CREATE TABLE src (c INT, d STRING) PARTITIONED BY (e INT);
    
    -- Add a partition to the src source table.
    ALTER TABLE src ADD IF NOT EXISTS PARTITION (e=201312);
    
    -- Append data to the src source table.
    INSERT INTO src PARTITION (e=201312) VALUES (1,100.1),(2,100.2),(3,100.3);
    
    -- Create a destination table named parttable.
    CREATE TABLE parttable(a INT, b DOUBLE) PARTITIONED BY (p STRING);
    
    -- Insert data from the src source table into the parttable destination table.
    SET odps.sql.allow.fullscan=true; 
    INSERT INTO parttable PARTITION (p) SELECT c, d, CURRENT_TIMESTAMP() FROM src;
    
    -- Query the parttable destination table.
    SET odps.sql.allow.fullscan=true;
    SELECT * FROM parttable;
    
    -- Output:
    +------------+------------+------------+
    | a          | b          | p          |
    +------------+------------+------------+
    | 1          | 100.1      | 2024-12-10 15:59:34.492 |
    | 2          | 100.2      | 2024-12-10 15:59:34.492 |
    | 3          | 100.3      | 2024-12-10 15:59:34.492 |
    +------------+------------+------------+
    Note

    Dynamic partition inserts can scatter ordered data, which may reduce the compression ratio. For better compression, use a Tunnel command to upload data to dynamic partitions. For a detailed example of how to use this command, see Migrate data from ApsaraDB RDS to MaxCompute based on dynamic partitioning.