Routine Load
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:
CSVorJSON. 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_fallbackparameter for the backend (BE) to the compatible legacy version. -
When you create a Routine Load, set
property.broker.version.fallbackto the older version that you want to be compatible with.
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 |
|
|
The name of the load job. Only one job with the same name can run at a time within the same database. |
|
|
Specifies the name of the destination table. |
|
|
Specifies the merge type. The default is |
|
|
Specifies data processing parameters for the load job. For more information, see load_properties. |
|
|
Specifies parameters for the load job. For more information, see job_properties Description. |
|
|
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 |
|
|
COLUMNS TERMINATED BY "," |
Specifies the column separator. Default: |
|
|
(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. |
|
|
None |
Specifies the conditions to filter source data. For more information, see data transformation. |
|
|
WHERE k1>100 and k2=1000 |
Specifies conditions to filter the imported data. For more information, see data transformation. |
|
|
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 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. |
|
|
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"
)
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" = "3" |
Specifies the desired concurrency. The value must be greater than 0. The default value is Note
|
|
|
"max_batch_interval" = "20" |
Specifies the maximum execution time for each subtask in seconds, with a default value of |
|
|
"max_batch_rows" = "300000" |
Specifies the maximum number of rows that each subtask can read. The default value is |
|
|
"max_batch_size" = "209715200" |
Specifies the maximum number of bytes that each subtask can read. The unit is bytes. The default value is |
|
|
"max_error_number"="3" |
Specifies the maximum number of error rows allowed in the specified sampling window. The default value is The sampling window is Note
Rows that are filtered out by the |
|
|
"strict_mode"="true" |
Specifies whether to enable strict mode. The default value is
|
|
|
"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" = "json" |
Specify the data import format. The default is |
|
|
-H "jsonpaths:[\"$.k2\",\"$.k1\"]" |
When the imported data format is |
|
|
-H "strip_outer_array:true" |
When you import data in |
|
|
-H "json_root:$.RECORDS" |
When you import data in JSON format, you can use |
|
|
— |
Specifies the parallelism for sending batch data. If the parallelism value exceeds the |
|
|
— |
Specifies whether to load data into only one tablet of the corresponding partition. The default is |
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 |
|
NULL |
|
Not null |
|
NULL |
true |
Invalid data (filtered) |
|
Not null |
aaa |
NULL |
false |
NULL |
|
Not null |
1 |
1 |
|
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 |
|
NULL |
|
Not null |
aaa |
NULL |
true |
Invalid data (filtered) |
|
Not null |
aaa |
NULL |
false |
NULL |
|
Not null |
|
1 |
|
Valid data |
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 |
|
|
Specifies the connection information for Kafka brokers. The format is Format: |
|
|
Specifies the Kafka topic to subscribe to. Format: |
|
|
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:
If you omit these parameters, the job subscribes to all topic partitions and starts consuming from Examples:
Important
You cannot mix timestamps and offsets in the same |
|
|
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 " |
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.protocolandproperty.ssl.ca.locationparameters 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 usekafka_default_offsetsto specify the starting offset. The default isOFFSET_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
-
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"); -
Import data by creating Routine Load jobs with different parameter settings:
-
Create a Kafka Routine Load job named
test1for thetest_tabletable in theexample_dbdatabase. 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
test2for thetest_tabletable in theexample_dbdatabase. 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_offsetto 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.
-
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; -
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 } ] -
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" );NoteThe partition key
dtin the table is not found in the example data, but is instead generated by thedt=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.
-
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 FILEcommand and set the catalog name tokafka.-
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"); -
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" );NoteSelectDB accesses Kafka clusters using the
librdkafkaC++ API. For parameters supported bylibrdkafka, see configuration properties.
-
-
Access a Kafka cluster with PLAIN authentication
To access a Kafka cluster that uses PLAIN authentication, add the following properties:
-
property.security.protocol=SASL_PLAINTEXT: Uses SASL plaintext. -
property.sasl.mechanism=PLAIN: Sets the SASL mechanism to PLAIN. -
property.sasl.username=admin: Sets the SASL username. -
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" ); -
-
Access a Kerberos-authenticated Kafka cluster
To access a Kerberos-authenticated Kafka cluster, add the following properties:
-
property.security.protocol=SASL_PLAINTEXT: Uses SASL plaintext. -
property.sasl.kerberos.service.name=$SERVICENAME: Specifies the Kerberos service name of the Kafka broker. -
property.sasl.kerberos.keytab=/etc/security/keytabs/${CLIENT_NAME}.keytab: Specifies the path to the client's Keytab file. -
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.conffile with the KDC service information. -
The value for
property.sasl.kerberos.keytabmust 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:
|
|
data_source |
The type of the data source. Currently, only |
|
data_source_properties |
The following properties are supported for the data source:
Note
Use |
Examples
-
The following example changes
desired_concurrent_numberto 1.ALTER ROUTINE LOAD FOR db1.label1 PROPERTIES ( "desired_concurrent_number" = "1" ); -
This example changes the
desired_concurrent_numberto 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. |
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_numThis 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_beThis 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_numAn 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_groupA 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_numAn 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_minAn 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_topicspecified 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 thenum.partitionsconfiguration 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=truein 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.listenersis configured in Kafka, the addresses specified inadvertised.listenersmust 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 thekafka_partitionslist. 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_ENDfor 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
STOPPEDstate, whereas jobs in thePAUSEDstate can be resumed.