Routine Load

Updated at:

Routine Load creates a persistent job that continuously reads data from a data source and loads it into ApsaraDB for SelectDB. This topic describes how to use Routine Load to load data from Kafka into an ApsaraDB for SelectDB instance.

Prerequisites

  • Supported data source: Currently, only Kafka is supported. You can connect to Kafka without authentication, or with PLAIN, SSL, or Kerberos authentication.

  • Supported message formats: CSV or JSON. For CSV, each message must be a single line without a trailing newline character.

  • Network connectivity: If your Kafka cluster is deployed on the Internet, the VPC in which your ApsaraDB for SelectDB instance resides must have outbound access to the Internet to connect to the Kafka broker. For more information about the configuration, see Resolve network issues with data sources.

Notes

By default, Kafka versions 0.10.0.0 and later are supported. If you use a Kafka version earlier than 0.10.0.0 (such as 0.9.0, 0.8.2, 0.8.1, or 0.8.0), you must enable backward compatibility:

  • Set the kafka_broker_version_fallback parameter for the backend (BE) to the compatible legacy version.

  • When you create a Routine Load, set property.broker.version.fallback to the older version that you want to be compatible with.

Note

Enabling backward compatibility may cause some newer Routine Load features to be unavailable. For example, you cannot set the offset of a Kafka partition by timestamp.

Import job

To use the Routine Load feature, you must first create a Routine Load job. This job continuously sends tasks through routine scheduling, and each task consumes a certain number of Kafka messages.

Syntax

CREATE ROUTINE LOAD [db.]job_name ON tbl_name
[merge_type]
[load_properties]
[job_properties]
FROM data_source [data_source_properties]

Parameters

Parameter

Description

[db.]job_name

The name of the load job. Only one job with the same name can run at a time within the same database.

tbl_name

Specifies the name of the destination table.

merge_type

Specifies the merge type. The default is APPEND, which indicates an append operation. The MERGE and DELETE types apply only to tables that use the unique key model. The MERGE type must be used with the [DELETE ON] clause to specify the delete flag column. The DELETE type indicates that all imported data represents deletions.

load_properties

Specifies data processing parameters for the load job. For more information, see load_properties.

job_properties

Specifies parameters for the load job. For more information, see job_properties Description.

data_source_properties

Specifies properties for the data source. For more information, see data_source_properties Description.

Load_properties

[column_separator],
[columns_mapping],
[preceding_filter],
[where_predicates],
[partitions],
[DELETE ON],
[ORDER BY]

Parameter

Example

Description

column_separator

COLUMNS TERMINATED BY ","

Specifies the column separator. Default: \t.

columns_mapping

(k1,k2,tmpk1,k3=tmpk1+1)

Specifies the mapping of source columns to target table columns, including various column transformations. For more information, see data transformation.

preceding_filter

None

Specifies the conditions to filter source data. For more information, see data transformation.

where_predicates

WHERE k1>100 and k2=1000

Specifies conditions to filter the imported data. For more information, see data transformation.

partitions

PARTITION(p1,p2,p3)

Specifies the target partitions to load data into. If not specified, the system automatically loads data into the corresponding partitions.

DELETE ON

DELETE ON v3>100

Specifies the column and expression that represent the delete flag.

Note

This parameter must be used with the MERGE import mode and applies only to unique key model tables.

ORDER BY

None

Specifies the sequence column to preserve row order during import.

Note

Applies only to tables that use the unique key model.

Job properties

PROPERTIES (
    "key1" = "val1",
    "key2" = "val2"
)
Note

The max_batch_interval, max_batch_rows, and max_batch_size parameters control the execution time and processing volume of each subtask. A subtask terminates when it reaches any of these thresholds.

Parameter

Example

Description

desired_concurrent_number

"desired_concurrent_number" = "3"

Specifies the desired concurrency. The value must be greater than 0. The default value is 3. A Routine Load job is divided into multiple tasks. This parameter specifies the maximum number of tasks that can run concurrently for a single job.

Note
  1. This is not the actual concurrency. The system determines the actual concurrency based on the number of nodes in the cluster, the current load, and the state of the data source.

  2. Increasing concurrency appropriately can accelerate processing by using a distributed cluster, but setting it too high can result in a large number of small file writes. The recommended value is Number of cluster cores / 16.

max_batch_interval

"max_batch_interval" = "20"

Specifies the maximum execution time for each subtask in seconds, with a default value of 10 and a valid range of 5 to 60 seconds.

max_batch_rows

"max_batch_rows" = "300000"

Specifies the maximum number of rows that each subtask can read. The default value is 200000. The value must be greater than or equal to 200000.

max_batch_size

"max_batch_size" = "209715200"

Specifies the maximum number of bytes that each subtask can read. The unit is bytes. The default value is 104857600 (100 MB). The valid values range from 100 MB to 1 GB.

max_error_number

"max_error_number"="3"

Specifies the maximum number of error rows allowed in the specified sampling window. The default value is 0, which means that no error rows are allowed. The value must be greater than or equal to 0.

The sampling window is max_batch_rows*10. This means that if the number of error rows within the sampling window exceeds a specified threshold, the routine job is paused, and manual intervention is required to check for data quality issues.

Note

Rows that are filtered out by the WHERE condition are not considered error rows.

strict_mode

"strict_mode"="true"

Specifies whether to enable strict mode. The default value is false. When strict mode is enabled, data is filtered during the import process if a column type conversion on non-empty source data results in a NULL value. The strict filtering policy is as follows:

  • During column type conversion, if strict mode is true, erroneous data is filtered. Erroneous data is any source data that is not a null value but becomes a null value after the conversion.

  • If an imported column is generated by a function transformation, strict mode does not apply.

  • For an imported column with a type that has a range constraint, if the source data can be successfully converted but falls outside the range, strict mode has no effect on the data. For example, if the column type is decimal(1,0) and the source data is 10, the data can be converted but is outside the range specified in the column declaration. Strict mode has no effect on this data.

timezone

"timezone" = "Africa/Abidjan"

Specifies the time zone for the import job. By default, the job uses the session's time zone.

Note

This parameter affects the results of all time zone-related functions used in the import.

format

"format" = "json"

Specify the data import format. The default is CSV, and the JSON format is also supported.

jsonpaths

-H "jsonpaths:[\"$.k2\",\"$.k1\"]"

When the imported data format is JSON, you can use jsonpaths to specify the fields to extract from the JSON data.

strip_outer_array

-H "strip_outer_array:true"

When you import data in JSON format, setting strip_outer_array to true specifies that the JSON data is an array, and each element in the array is treated as a single row of data. The default value is false.

json_root

-H "json_root:$.RECORDS"

When you import data in JSON format, you can use json_root to specify the root node of the JSON data. SelectDB uses json_root to extract and parse the elements of the root node. The default value is empty.

send_batch_parallelism

—

Specifies the parallelism for sending batch data. If the parallelism value exceeds the max_send_batch_parallelism_per_job value in the BE configuration, the coordinator node (BE) uses the max_send_batch_parallelism_per_job value.

load_to_single_tablet

—

Specifies whether to load data into only one tablet of the corresponding partition. The default is false. You can set this parameter only when importing data into a duplicate table that uses random partitioning.

Effect of strict_mode on source data import

The following example is for a column of the TinyInt type where the column allows null values.

Source data

Source data example

String to int

Strict_mode

Result

Null value

\N

N/A

true or false

NULL

Not null

aaa or 2000

NULL

true

Invalid data (filtered)

Not null

aaa

NULL

false

NULL

Not null

1

1

true or false

Valid data

The following example is for a column of the Decimal(1,0) type where the column allows null values.

Source data

Source data example

String to int

Strict_mode

Result

Null value

\N

N/A

true or false

NULL

Not null

aaa

NULL

true

Invalid data (filtered)

Not null

aaa

NULL

false

NULL

Not null

1 or 10

1

true or false

Valid data

Note

Although 10 is out of range, it conforms to the Decimal type, so strict_mode does not affect it. A later ETL processing step filters out the value 10, but strict_mode does not.

data_source_properties parameters

FROM KAFKA
(
    "key1" = "val1",
    "key2" = "val2"
)

Parameter

Description

kafka_broker_list

Specifies the connection information for Kafka brokers. The format is ip:host. Separate multiple brokers with a comma.

Format: "kafka_broker_list"="broker1:9092,broker2:9092".

kafka_topic

Specifies the Kafka topic to subscribe to.

Format: "kafka_topic"="my_topic".

kafka_partitions/kafka_offsets

Specifies the Kafka partitions to subscribe to and the corresponding starting offset for each partition. If you specify a timestamp, consumption starts from the first offset recorded at or after that time.

An offset can be an integer greater than or equal to 0, or one of the following values:

  • OFFSET_BEGINNING: Starts consumption from the earliest available data.

  • OFFSET_END: Starts consumption from the latest offset.

  • A timestamp, for example: "2021-05-22 11:00:00"

If you omit these parameters, the job subscribes to all topic partitions and starts consuming from OFFSET_END by default.

Examples:

"kafka_partitions" = "0,1,2,3",
"kafka_offsets" = "101,0,OFFSET_BEGINNING,OFFSET_END"
"kafka_partitions" = "0,1,2,3",
"kafka_offsets" = "2021-05-22 11:00:00,2021-05-22 11:00:00,2021-05-22 11:00:00"
Important

You cannot mix timestamps and offsets in the same kafka_offsets list.

property

Specifies custom Kafka parameters. This is equivalent to the "--property" parameter in the Kafka shell.

For file path values, prefix the value with the keyword "FILE:".

Property parameters

  • To connect to Kafka using SSL, specify the following parameters:

    "property.security.protocol" = "ssl",
    "property.ssl.ca.location" = "FILE:ca.pem",
    "property.ssl.certificate.location" = "FILE:client.pem",
    "property.ssl.key.location" = "FILE:client.key",
    "property.ssl.key.password" = "abcdefg"

    The property.security.protocol and property.ssl.ca.location parameters are required. These parameters specify an SSL connection and provide the location of the CA certificate.

    If client authentication is enabled on the Kafka server, you must also set the following parameters:

    "property.ssl.certificate.location"
    "property.ssl.key.location"
    "property.ssl.key.password"

    These parameters specify the client's public key, private key, and the password for the private key, respectively.

  • Specifies the default starting offset for Kafka partitions.

    If you do not specify kafka_partitions/kafka_offsets, the job consumes all partitions by default. In this case, you can use kafka_default_offsets to specify the starting offset. The default is OFFSET_END, which means consumption starts from the latest offset.

    "property.kafka_default_offsets" = "OFFSET_BEGINNING"

For additional supported custom parameters, see the client-side configuration options in the official librdkafka CONFIGURATION documentation. For example:

"property.client.id" = "12345",
"property.ssl.ca.location" = "FILE:ca.pem"

Usage example

Simple Routine Load job

  1. Create the destination SelectDB table for data import:

    CREATE TABLE test_table
    (
        id int,
        name varchar(50),
        age int,
        address varchar(50),
        url varchar(500)
    )
    UNIQUE KEY(`id`, `name`)
    DISTRIBUTED BY HASH(id) BUCKETS 4
    PROPERTIES("replication_num" = "1");
  2. Import data by creating Routine Load jobs with different parameter settings:

    • Create a Kafka Routine Load job named test1 for the test_table table in the example_db database. This job uses a comma as the column delimiter, consumes data from all partitions, and starts from the earliest available offset (OFFSET_BEGINNING).

      CREATE ROUTINE LOAD example_db.test1 ON test_table
      COLUMNS TERMINATED BY ",",
      COLUMNS(k1, k2, k3, v1, v2, v3 = k1 * 100)
      PROPERTIES
      (
          "desired_concurrent_number"="3",
          "max_batch_interval" = "20",
          "max_batch_rows" = "300000",
          "max_batch_size" = "209715200",
          "strict_mode" = "false"
      )
      FROM KAFKA
      (
          "kafka_broker_list" = "broker1:9092,broker2:9092,broker3:9092",
          "kafka_topic" = "my_topic",
          "property.kafka_default_offsets" = "OFFSET_BEGINNING"
      );
    • Create a Kafka Routine Load job named test2 for the test_table table in the example_db database. This example enables strict mode.

      CREATE ROUTINE LOAD example_db.test2 ON test_table
      COLUMNS TERMINATED BY ",",
      COLUMNS(k1, k2, k3, v1, v2, v3 = k1 * 100)
      PROPERTIES
      (
          "desired_concurrent_number"="3",
          "max_batch_interval" = "20",
          "max_batch_rows" = "300000",
          "max_batch_size" = "209715200",
          "strict_mode" = "true"
      )
      FROM KAFKA
      (
          "kafka_broker_list" = "broker1:9092,broker2:9092,broker3:9092",
          "kafka_topic" = "my_topic",
          "property.kafka_default_offsets" = "OFFSET_BEGINNING"
      );
    • To start consuming data from a specific point in time, set property.kafka_default_offset to a timestamp.

      CREATE ROUTINE LOAD example_db.test4 ON test_table
      PROPERTIES
      (
          "desired_concurrent_number"="3",
          "max_batch_interval" = "30",
          "max_batch_rows" = "300000",
          "max_batch_size" = "209715200"
      ) FROM KAFKA
      (
          "kafka_broker_list" = "broker1:9092,broker2:9092",
          "kafka_topic" = "my_topic",
          "property.kafka_default_offset" = "2024-01-21 10:00:00"
      );

Load JSON data

Routine Load supports only the following two JSON formats.

  • A single record formatted as a JSON object.

    For a single-table import (specified with ON TABLE_NAME), use the following JSON data format:

    {"key1":"value1","key2":"value2","key3":"value3"}

    For a dynamic or multi-table Routine Load import (without a specified table name), use the following JSON data format:

    table_name|{"key1":"value1","key2":"value2","key3":"value3"}
  • A JSON array of multiple records.

    For a single-table import (specified with ON TABLE_NAME), use the following JSON data format:

    [
        {   
            "key1":"value11",
            "key2":"value12",
            "key3":"value13",
            "key4":14
        },
        {
            "key1":"value21",
            "key2":"value22",
            "key3":"value23",
            "key4":24
        },
        {
            "key1":"value31",
            "key2":"value32",
            "key3":"value33",
            "key4":34
        }
    ]

    For a dynamic or multi-table Routine Load import (without a specified table name), use the following JSON data format:

       table_name|[
        {   
            "key1":"value11",
            "key2":"value12",
            "key3":"value13",
            "key4":14
        },
        {
            "key1":"value21",
            "key2":"value22",
            "key3":"value23",
            "key4":24
        },
        {
            "key1":"value31",
            "key2":"value32",
            "key3":"value33",
            "key4":34
        }
    ]

The following example shows how to import data in JSON format.

  1. Create the SelectDB table for the import as follows:

    CREATE TABLE `example_tbl` (
       `category` varchar(24) NULL COMMENT "",
       `author` varchar(24) NULL COMMENT "",
       `timestamp` bigint(20) NULL COMMENT "",
       `dt` int(11) NULL COMMENT "",
       `price` double REPLACE
    ) ENGINE=OLAP
    AGGREGATE KEY(`category`,`author`,`timestamp`,`dt`)
    COMMENT "OLAP"
    PARTITION BY RANGE(`dt`)
    (
      PARTITION p0 VALUES [("-2147483648"), ("20230509")),
        PARTITION p20200509 VALUES [("20230509"), ("20231010")),
        PARTITION p20200510 VALUES [("20231010"), ("20231211")),
        PARTITION p20200511 VALUES [("20231211"), ("20240512"))
    )
    DISTRIBUTED BY HASH(`category`,`author`,`timestamp`) BUCKETS 4;
  2. The Kafka topic contains two types of JSON records:

    {
        "category":"value1331",
        "author":"value1233",
        "timestamp":1700346050,
        "price":1413
    }
    [
        {
            "category":"value13z2",
            "author":"vaelue13",
            "timestamp":1705645251,
            "price":14330
        },
        {
            "category":"lvalue211",
            "author":"lvalue122",
            "timestamp":1684448450,
            "price":24440
        }
    ]
  3. The following examples show how to load JSON data in different modes.

    • Loading JSON data in simple mode.

      CREATE ROUTINE LOAD example_db.test_json_label_1 ON example_tbl
      COLUMNS(category,price,author)
      PROPERTIES
      (
          "desired_concurrent_number"="3",
          "max_batch_interval" = "20",
          "max_batch_rows" = "300000",
          "max_batch_size" = "209715200",
          "strict_mode" = "false",
          "format" = "json"
      )
      FROM KAFKA
      (
          "kafka_broker_list" = "broker1:9092,broker2:9092,broker3:9092",
          "kafka_topic" = "my_topic",
          "kafka_partitions" = "0,1,2",
          "kafka_offsets" = "0,0,0"
       );
    • Precisely importing JSON data.

      CREATE ROUTINE LOAD example_db.test_json_label_3 ON example_tbl
      COLUMNS(category, author, price, timestamp, dt=from_unixtime(timestamp, '%Y%m%d'))
      PROPERTIES
      (
          "desired_concurrent_number"="3",
          "max_batch_interval" = "20",
          "max_batch_rows" = "300000",
          "max_batch_size" = "209715200",
          "strict_mode" = "false",
          "format" = "json",
          "jsonpaths" = "[\"$.category\",\"$.author\",\"$.price\",\"$.timestamp\"]",
          "strip_outer_array" = "true"
      )
      FROM KAFKA
      (
          "kafka_broker_list" = "broker1:9092,broker2:9092,broker3:9092",
          "kafka_topic" = "my_topic",
          "kafka_partitions" = "0,1,2",
          "kafka_offsets" = "0,0,0"
      );
      Note

      The partition key dt in the table is not found in the example data, but is instead generated by the dt=from_unixtime(timestamp,'%Y%m%d') expression in the Routine Load statement.

Kafka authentication methods

The following examples illustrate how to access a Kafka cluster with different authentication methods.

  1. Access an SSL-authenticated Kafka cluster

    To access an SSL-authenticated Kafka cluster, provide a certificate file (ca.pem) to authenticate the Kafka broker's public key. If the cluster also requires client authentication, provide the client's public key (client.pem), private key file (client.key), and key password. First, upload these files to SelectDB using the CREATE FILE command and set the catalog name to kafka.

    1. Upload the files. Example:

      CREATE FILE "ca.pem" PROPERTIES("url" = "https://example_url/kafka-key/ca.pem", "catalog" = "kafka");
      CREATE FILE "client.key" PROPERTIES("url" = "https://example_urlkafka-key/client.key", "catalog" = "kafka");
      CREATE FILE "client.pem" PROPERTIES("url" = "https://example_url/kafka-key/client.pem", "catalog" = "kafka");
    2. Create a Routine Load job. Example:

      CREATE ROUTINE LOAD db1.job1 on tbl1
      PROPERTIES
      (
          "desired_concurrent_number"="1"
      )
      FROM KAFKA
      (
          "kafka_broker_list"= "broker1:9091,broker2:9091",
          "kafka_topic" = "my_topic",
          "property.security.protocol" = "ssl",
          "property.ssl.ca.location" = "FILE:ca.pem",
          "property.ssl.certificate.location" = "FILE:client.pem",
          "property.ssl.key.location" = "FILE:client.key",
          "property.ssl.key.password" = "abcdefg"
      );
      Note

      SelectDB accesses Kafka clusters using the librdkafka C++ API. For parameters supported by librdkafka, see configuration properties.

  2. Access a Kafka cluster with PLAIN authentication

    To access a Kafka cluster that uses PLAIN authentication, add the following properties:

    1. property.security.protocol=SASL_PLAINTEXT: Uses SASL plaintext.

    2. property.sasl.mechanism=PLAIN: Sets the SASL mechanism to PLAIN.

    3. property.sasl.username=admin: Sets the SASL username.

    4. property.sasl.password=admin: Sets the SASL password.

    Create a Routine Load job. Example:

    CREATE ROUTINE LOAD db1.job1 on tbl1
    PROPERTIES (
    "desired_concurrent_number"="1",
     )
    FROM KAFKA
    (
        "kafka_broker_list" = "broker1:9092,broker2:9092",
        "kafka_topic" = "my_topic",
        "property.security.protocol"="SASL_PLAINTEXT",
        "property.sasl.mechanism"="PLAIN",
        "property.sasl.username"="admin",
        "property.sasl.password"="admin"
    );
    
  3. Access a Kerberos-authenticated Kafka cluster

    To access a Kerberos-authenticated Kafka cluster, add the following properties:

    1. property.security.protocol=SASL_PLAINTEXT: Uses SASL plaintext.

    2. property.sasl.kerberos.service.name=$SERVICENAME: Specifies the Kerberos service name of the Kafka broker.

    3. property.sasl.kerberos.keytab=/etc/security/keytabs/${CLIENT_NAME}.keytab: Specifies the path to the client's Keytab file.

    4. property.sasl.kerberos.principal=${CLIENT_NAME}/${CLIENT_HOST}: Specifies the Kerberos principal for the SelectDB client.

    Create a Routine Load job. Example:

    CREATE ROUTINE LOAD db1.job1 on tbl1
    PROPERTIES (
    "desired_concurrent_number"="1",
     )
    FROM KAFKA
    (
        "kafka_broker_list" = "broker1:9092,broker2:9092",
        "kafka_topic" = "my_topic",
        "property.security.protocol" = "SASL_PLAINTEXT",
        "property.sasl.kerberos.service.name" = "kafka",
        "property.sasl.kerberos.keytab" = "/etc/krb5.keytab",
        "property.sasl.kerberos.principal" = "id@your.com"
    );
    Note
    • To allow SelectDB to access a Kerberos-authenticated Kafka cluster, install the Kerberos client (kinit) on all running nodes in the SelectDB cluster. Also, configure the krb5.conf file with the KDC service information.

    • The value for property.sasl.kerberos.keytab must be the absolute path to the Keytab file, and the SelectDB process must have read permission for this file.

Modify a Routine Load job

Modifies an existing Routine Load job. You can only modify jobs in the PAUSED state.

Syntax

ALTER ROUTINE LOAD FOR <job_name>
[job_properties]
FROM <data_source>
[data_source_properties]

Parameters

Parameter

Description

[db.]job_name

The name of the job to modify.

tbl_name

The name of the destination table for the load.

job_properties

The following job properties can be modified:

  • desired_concurrent_number

  • max_error_number

  • max_batch_interval

  • max_batch_rows

  • max_batch_size

  • jsonpaths

  • json_root

  • strip_outer_array

  • strict_mode

  • timezone

  • num_as_string

  • fuzzy_parse

data_source

The type of the data source. Currently, only Kafka is supported.

data_source_properties

The following properties are supported for the data source:

  1. kafka_partitions

  2. kafka_offsets

  3. kafka_broker_list

  4. kafka_topic

  5. Custom properties, such as property.group.id.

Note

Use kafka_partitions and kafka_offsets to change the consumer offsets for specific Kafka partitions. You can only change offsets for partitions already being consumed, and you cannot add new partitions.

Examples

  • The following example changes desired_concurrent_number to 1.

    ALTER ROUTINE LOAD FOR db1.label1
    PROPERTIES
    (
        "desired_concurrent_number" = "1"
    );
  • This example changes the desired_concurrent_number to 10, the partition offsets, and the group ID.

    ALTER ROUTINE LOAD FOR db1.label1
    PROPERTIES
    (
        "desired_concurrent_number" = "10"
    )
    FROM kafka
    (
        "kafka_partitions" = "0, 1, 2",
        "kafka_offsets" = "100, 200, 100",
        "property.group.id" = "new_group"
    );

Pause Routine Load job

This command pauses a Routine Load job. To resume the job, use the RESUME command.

Syntax

PAUSE [ALL] ROUTINE LOAD FOR <job_name>;

Parameters

Parameter

Description

[db.]job_name

The name of the job to pause.

Examples

  • Pause the Routine Load job named test1:

    PAUSE ROUTINE LOAD FOR test1;
  • Pause all Routine Load jobs:

    PAUSE ALL ROUTINE LOAD;

Resume a Routine Load job

Resumes a paused Routine Load job. The job continues to consume from its last recorded offset.

Syntax

RESUME [ALL] ROUTINE LOAD FOR <job_name>

Parameters

Parameter

Description

[db.]job_name

The name of the job to resume.

Examples

  • Resume the Routine Load job named test1:

    RESUME ROUTINE LOAD FOR test1;
  • Resume all Routine Load jobs:

    RESUME ALL ROUTINE LOAD;

Stop a Routine Load job

Stops a Routine Load job. A stopped job cannot be restarted. This action does not roll back any imported data.

Syntax

STOP ROUTINE LOAD FOR <job_name>;

Parameters

Parameter

Description

[db.]job_name

The name of the Routine Load job to stop.

Example

Stop the Routine Load job named test1:

STOP ROUTINE LOAD FOR test1;

View import jobs

Use the SHOW ROUTINE LOAD command to check the status of a Routine Load job.

Syntax

SHOW [ALL] ROUTINE LOAD [FOR job_name];

Parameters

Parameter

Description

[db.]job_name

Specifies the name of the job to view.

Note

If the data format is invalid, detailed error messages are recorded in the ErrorLogUrls field, which may contain multiple links. To view the errors, copy a link and open it in your browser.

Examples

  • Display all Routine Load jobs named test1, including those that are stopped or canceled. The result may contain one or more rows.

    SHOW ALL ROUTINE LOAD FOR test1;
  • Display the currently running Routine Load job named test1.

    SHOW ROUTINE LOAD FOR test1;
  • Display all Routine Load jobs in the example_db database, including those that are stopped or canceled. The result may contain one or more rows.

    use example_db;
    SHOW ALL ROUTINE LOAD;
  • Display all currently running Routine Load jobs in the example_db database.

    use example_db;
    SHOW ROUTINE LOAD;
  • Display the currently running Routine Load job named test1 in the example_db database.

    SHOW ROUTINE LOAD FOR example_db.test1;
  • Display all Routine Load jobs named test1 in the example_db database, including those that are stopped or canceled. The result may contain one or more rows.

    SHOW ALL ROUTINE LOAD FOR example_db.test1;

System configuration​

The following system configuration parameters affect Routine Load.

  • max_routine_load_task_concurrent_num

    This FE parameter limits the maximum number of concurrent tasks within a single Routine Load job. It defaults to 5 and can be modified at runtime. We recommend keeping the default value, as a higher value can exhaust cluster resources.

  • max_routine_load_task_num_per_be

    This FE parameter limits the maximum number of concurrent tasks on each BE node. It defaults to 5 and can be modified at runtime. We recommend keeping the default value, as a higher value can exhaust cluster resources.

  • max_routine_load_job_num

    An FE parameter with a default value of 100. This parameter can be modified at runtime. It limits the total number of Routine Load jobs in the NEED_SCHEDULED, RUNNING, or PAUSED states. New job submissions are blocked once this limit is reached.

  • max_consumer_num_per_group

    A BE parameter with a default value of 3. This parameter specifies the maximum number of consumers that can be created per task. For a Kafka data source, one consumer can handle one or more Kafka partitions. For example, if a task needs to consume 6 Kafka partitions, 3 consumers are created, each handling 2 partitions. If there are only 2 partitions, 2 consumers are created, each handling 1 partition.

  • max_tolerable_backend_down_num

    An FE parameter with a default value of 0. Under certain conditions, ApsaraDB for SelectDB can automatically resume PAUSED tasks by changing their state to RUNNING. A value of 0 means that a task can be automatically resumed only when all BE nodes are available.

  • period_of_auto_resume_min

    An FE parameter with a default value of 5 (in minutes). When ApsaraDB for SelectDB automatically resumes a task, it makes up to three attempts within this period. If all attempts fail, the task is locked, will not be rescheduled automatically, and requires manual recovery.

Considerations​

  • Relationship between Routine Load jobs and ALTER TABLE operations

    • Routine Load jobs do not block schema change or ROLLUP operations. However, after a schema change, an invalid column mapping can cause a spike in error rows, eventually causing the job to enter the PAUSED state. To mitigate this, we recommend explicitly specifying the column mapping in your Routine Load job and adding nullable columns or columns with default values.

    • Deleting a table partition may prevent an import job from finding the corresponding partition, causing the job to enter the PAUSED state.

  • Relationship between Routine Load jobs and other import jobs (LOAD, DELETE, INSERT)

    • Routine Load jobs do not conflict with other LOAD jobs or INSERT operations.

    • When you perform a DELETE operation, the target table partition cannot have any running import tasks. Therefore, before executing DELETE, you must pause the Routine Load job and wait for all dispatched tasks to complete.

  • Relationship between Routine Load jobs and DROP DATABASE or DROP TABLE operations

    When the database or table associated with a Routine Load job is dropped, the job is automatically CANCELED.

  • Relationship between Kafka-based Routine Load jobs and Kafka topics

    If the kafka_topic specified for a Routine Load job does not exist in the Kafka cluster:

    • If the brokers in your Kafka cluster are set to auto.create.topics.enable=true, topics are automatically created. The number of partitions for an automatically created topic is determined by the num.partitions configuration of the brokers in your Kafka cluster. A routine job will continuously read data from the topic.

    • If a Kafka broker is configured with auto.create.topics.enable=false, the topic is not created. The job enters the PAUSED state before reading any data.

    Therefore, to have Kafka automatically create a topic for your Routine Load job if the topic does not exist, set auto.create.topics.enable=true in the broker configuration.

  • If your environment has network segment isolation or custom DNS resolution, ensure the following:

    • The broker list specified for the Routine Load job must be accessible from the ApsaraDB for SelectDB service.

    • If advertised.listeners is configured in Kafka, the addresses specified in advertised.listeners must also be accessible from the ApsaraDB for SelectDB service.

  • Specifying partitions and offsets for consumption

    ApsaraDB for SelectDB allows you to specify partitions and offsets for consumption. These parameters can be used in combination.

    The three related parameters are as follows:

    • kafka_partitions: Specifies the list of partitions to consume from, such as "0,1,2,3".

    • kafka_offsets: Specifies the starting offset for each partition. The number of offsets must match the number of partitions in the kafka_partitions list. For example, "1000,1000,2000,2000".

    • property.kafka_default_offset: Specifies the default starting offset for partitions.

    When creating an import job, these three parameters can be combined in the following five ways:

    Case

    kafka_partitions

    kafka_offsets

    property.kafka_default_offset

    Behavior

    1

    No

    No

    No

    The system automatically discovers all partitions for the topic and starts consuming from OFFSET_END.

    2

    No

    No

    Yes

    The system automatically discovers all partitions for the topic and starts consuming from the specified default offset.

    3

    Yes

    No

    No

    The system consumes from the specified partitions, starting from OFFSET_END for each.

    4

    Yes

    Yes

    No

    The system consumes from the specified partitions, starting from the specified offset for each.

    5

    Yes

    No

    Yes

    The system consumes from the specified partitions, starting from the specified default offset for each.

  • Difference between STOPPED and PAUSED

    The FE periodically purges jobs in the STOPPED state, whereas jobs in the PAUSED state can be resumed.