Import data by using Spark
ApsaraDB for SelectDB is compatible with Apache Doris and lets you use the Spark Doris Connector to import large volumes of data using Spark's distributed computing power. This topic explains how the connector works and shows how to use it to synchronize data to ApsaraDB for SelectDB.
Overview
The Spark Doris Connector is one of the methods for importing large amounts of data into Alibaba Cloud SelectDB. Using Spark's distributed computing capabilities, you can read large datasets from various upstream data sources, such as MySQL, PostgreSQL, HDFS, and S3, into a DataFrame. You can then use the Spark Doris Connector to load the data into an ApsaraDB for SelectDB table. You can also use Spark's JDBC interface to read data from tables in ApsaraDB for SelectDB.
How it works
The following figure shows the architecture for importing data into ApsaraDB for SelectDB by using the Spark Doris Connector. In this architecture, the Spark Doris Connector acts as a bridge for writing external data to ApsaraDB for SelectDB. It uses Spark's distributed computing cluster to preprocess data, which accelerates the entire data pipeline. This high-performance approach replaces traditional data ingestion over a JDBC connection.
Prerequisites
To import data by using the Spark Doris Connector, you must use version 1.3.1 or later of the connector package.
Add the Spark Doris Connector dependency
Add the Doris Connector dependency in one of the following ways:
If you use Maven, add the dependency as shown in the following code. For other versions, see Maven Repository.
<dependency> <groupId>org.apache.doris</groupId> <artifactId>spark-doris-connector-3.2_2.12</artifactId> <version>1.3.2</version> </dependency>Add the connector by using its JAR package.
The following table lists three common connectors. Select a connector package that matches your Spark version. For more versions, see Maven Repository.
NoteThe following JAR packages are compiled with Java 8. If you require a different Java version, contact ApsaraDB for SelectDB technical support.
In the following table, the Connector column indicates the supported Spark, Scala, and connector versions, respectively.
Connector
Runtime JAR
2.4-2.12-1.3.2
3.1-2.12-1.3.2
3.2-2.12-1.3.2
After you obtain the JAR package, use it in one of the following ways:
If you run Spark in local mode, place the downloaded JAR package in the
jarsdirectory of your Spark installation.If you run Spark in YARN cluster mode, add the JAR package to your pre-deployment package. For example:
Upload the
spark-doris-connector-3.2_2.12-1.3.2.jarfile to HDFS.hdfs dfs -mkdir /spark-jars/ hdfs dfs -put /<your_local_path>/spark-doris-connector-3.2_2.12-1.3.2.jar/spark-jars/Add the
spark-doris-connector-3.2_2.12-1.3.2.jardependency to the cluster.spark.yarn.jars=hdfs:///spark-jars/spark-doris-connector-3.2_2.12-1.3.2.jar
Usage
After you run Spark on a Spark client or add the connector package to your Spark development environment, you can synchronize data by using either Spark SQL or the DataFrame API. The following examples show how to synchronize data from Spark to ApsaraDB for SelectDB.
Spark SQL
val selectdbHttpPort = "selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080"
val selectdbJdbc = "jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030"
val selectdbUser = "admin"
val selectdbPwd = "****"
val selectdbTable = "test_db.test_order"
CREATE TEMPORARY VIEW test_order
USING doris
OPTIONS(
"table.identifier"="${selectdbTable}",
"fenodes"="${selectdbHttpPort}",
"user"="${selectdbUser}",
"password"="${selectdbPwd}",
"sink.properties.format"="json"
);
INSERT INTO test_order SELECT order_id,order_amount,order_status FROM tmp_tb;Parameters
Parameter | Default | Required | Description |
fenodes | None | Yes | The HTTP endpoint of your ApsaraDB for SelectDB instance. You can get the VPC Endpoint (or Public Endpoint) and HTTP Port from the Instance Details > Network Information page in the ApsaraDB for SelectDB console. Example: |
table.identifier | None | Yes | The destination table in your ApsaraDB for SelectDB instance. Format: |
request.retries | 3 | No | The number of retries for requests to SelectDB. |
request.connect.timeout.ms | 30000 | No | The connection timeout for requests to SelectDB, in milliseconds. |
request.read.timeout.ms | 30000 | No | The read timeout for requests to SelectDB, in milliseconds. |
request.query.timeout.s | 3600 | No | The query timeout for SelectDB, in seconds. The default value is 1 hour. A value of |
request.tablet.size | Integer.MAX_VALUE | No | The number of SelectDB tablets per Spark RDD partition. A smaller value generates more partitions. This increases parallelism on the Spark side but also increases the load on SelectDB. |
read.field | None | No | A comma-separated list of column names to read from the SelectDB table. |
batch.size | 1024 | No | The maximum number of rows to read from a BE in a single batch. Increasing this value can reduce the number of connections between Spark and SelectDB, minimizing network latency overhead. |
exec.mem.limit | 2147483648 | No | The memory limit for a single query. The default is 2 GB. Unit: bytes. |
deserialize.arrow.async | false | No | Specifies whether to asynchronously deserialize Arrow format data to the RowBatch required by the spark-doris-connector iterator. |
deserialize.queue.size | 64 | No | The internal processing queue size for asynchronous Arrow deserialization. This parameter takes effect only when |
write.fields | None | No | Specifies the fields or the order of fields to write to the SelectDB table. Use commas to separate multiple columns. By default, all fields are written in the order they are defined in the SelectDB table. |
sink.batch.size | 100000 | No | The maximum number of rows to write to a BE in a single request. |
sink.max-retries | 0 | No | The number of retries after a failed write attempt to a BE. |
sink.properties.format | csv | No | The data format for Stream Load. Valid values: |
sink.properties.* | -- | No | Parameters for the Stream Load job. For example, to specify a column separator, use |
sink.task.partition.size | None | No | The number of partitions for the SelectDB write task. After operations like filtering, a Spark RDD might have many partitions, each containing only a few records. This increases write frequency and wastes computing resources. A smaller value can reduce the write frequency and decrease the compaction pressure on SelectDB. This parameter is used with |
sink.task.use.repartition | false | No | Specifies whether to use repartition to control the number of write partitions for SelectDB. If If |
sink.batch.interval.ms | 50 | No | The interval between each sink batch, in milliseconds. |
sink.enable-2pc | false | No | Specifies whether to enable two-phase commit (2PC). If enabled, transactions are committed only when the job completes. If any task fails, all pre-committed transactions are rolled back. |
sink.auto-redirect | true | No | Specifies whether to redirect Stream Load requests. If enabled, Stream Load writes data through the FE, and you no longer need to explicitly get BE information. |
user | None | Yes | The username to access your ApsaraDB for SelectDB instance. |
password | None | Yes | The password to access your ApsaraDB for SelectDB instance. |
filter.query.in.max.count | 100 | No | The maximum number of elements in an |
ignore-type | None | No | A comma-separated list of field types to ignore when reading the schema from a temporary view. Example: |
DataFrame
val spark = SparkSession.builder().master("local[1]").getOrCreate()
val df = spark.createDataFrame(Seq(
("1", 100, "Pending payment"),
("2", 200, null),
("3", 300, "Received")
)).toDF("order_id", "order_amount", "order_status")
df.write
.format("doris")
.option("fenodes", selectdbHttpPort)
.option("table.identifier", selectdbTable)
.option("user", selectdbUser)
.option("password", selectdbPwd)
.option("sink.batch.size", 100000)
.option("sink.max-retries", 3)
.option("sink.properties.file.column_separator", "\t")
.option("sink.properties.file.line_delimiter", "\n")
.save()Parameters
Parameter | Default | Required | Description |
fenodes | None | Yes | The HTTP endpoint of your ApsaraDB for SelectDB instance. You can get the VPC Endpoint (or Public Endpoint) and HTTP Port from the Instance Details > Network Information page in the ApsaraDB for SelectDB console. Example: |
table.identifier | None | Yes | The destination table in your ApsaraDB for SelectDB instance. Format: |
request.retries | 3 | No | The number of retries for requests to SelectDB. |
request.connect.timeout.ms | 30000 | No | The connection timeout for requests to SelectDB, in milliseconds. |
request.read.timeout.ms | 30000 | No | The read timeout for requests to SelectDB, in milliseconds. |
request.query.timeout.s | 3600 | No | The query timeout for SelectDB, in seconds. The default value is 1 hour. A value of |
request.tablet.size | Integer.MAX_VALUE | No | The number of SelectDB tablets per Spark RDD partition. A smaller value generates more partitions. This increases parallelism on the Spark side but also increases the load on SelectDB. |
read.field | None | No | A comma-separated list of column names to read from the SelectDB table. |
batch.size | 1024 | No | The maximum number of rows to read from a BE in a single batch. Increasing this value can reduce the number of connections between Spark and SelectDB, minimizing network latency overhead. |
exec.mem.limit | 2147483648 | No | The memory limit for a single query. The default is 2 GB. Unit: bytes. |
deserialize.arrow.async | false | No | Specifies whether to asynchronously deserialize Arrow format data to the RowBatch required by the spark-doris-connector iterator. |
deserialize.queue.size | 64 | No | The internal processing queue size for asynchronous Arrow deserialization. This parameter takes effect only when |
write.fields | None | No | Specifies the fields or the order of fields to write to the SelectDB table. Use commas to separate multiple columns. By default, all fields are written in the order they are defined in the SelectDB table. |
sink.batch.size | 100000 | No | The maximum number of rows to write to a BE in a single request. |
sink.max-retries | 0 | No | The number of retries after a failed write attempt to a BE. |
sink.properties.format | csv | No | The data format for Stream Load. Valid values: |
sink.properties.* | -- | No | Parameters for the Stream Load job. For example, to specify a column separator, use |
sink.task.partition.size | None | No | The number of partitions for the SelectDB write task. After operations like filtering, a Spark RDD might have many partitions, each containing only a few records. This increases write frequency and wastes computing resources. A smaller value can reduce the write frequency and decrease the compaction pressure on SelectDB. This parameter is used with |
sink.task.use.repartition | false | No | Specifies whether to use repartition to control the number of write partitions for SelectDB. If If |
sink.batch.interval.ms | 50 | No | The interval between each sink batch, in milliseconds. |
sink.enable-2pc | false | No | Specifies whether to enable two-phase commit (2PC). If enabled, transactions are committed only when the job completes. If any task fails, all pre-committed transactions are rolled back. |
sink.auto-redirect | true | No | Specifies whether to redirect Stream Load requests. If enabled, Stream Load writes data through the FE, and you no longer need to explicitly get BE information. |
user | None | Yes | The username to access your ApsaraDB for SelectDB instance. |
password | None | Yes | The password to access your ApsaraDB for SelectDB instance. |
filter.query.in.max.count | 100 | No | The maximum number of elements in an |
ignore-type | None | No | A comma-separated list of field types to ignore when reading the schema from a temporary view. Example: |
sink.streaming.passthrough | false | No | Writes the values in the first column without processing. |
Examples
The following table lists the software versions used in the example environment.
Software | Java | Spark | Scala | SelectDB |
Version | 1.8 | 3.1.2 | 2.12 | 3.0.4 |
Environment preparation
Configure the Spark environment.
Download and decompress the Spark installation package. This example uses the
spark-3.1.2-bin-hadoop3.2.tgzpackage.wget https://archive.apache.org/dist/spark/spark-3.1.2/spark-3.1.2-bin-hadoop3.2.tgz tar xvzf spark-3.1.2-bin-hadoop3.2.tgzPlace the
spark-doris-connector-3.2_2.12-1.3.2.jarfile in theSPARK_HOME/jarsdirectory.
Prepare the data for import. This example uses a small amount of sample data from a MySQL database.
Create a test table in MySQL.
CREATE TABLE `employees` ( `emp_no` int NOT NULL, `birth_date` date NOT NULL, `first_name` varchar(14) NOT NULL, `last_name` varchar(16) NOT NULL, `gender` enum('M','F') NOT NULL, `hire_date` date NOT NULL, PRIMARY KEY (`emp_no`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb3Use DMS to generate test data. For more information, see Generate test data.
Configure the ApsaraDB for SelectDB instance.
Create an ApsaraDB for SelectDB instance. For more information, see Create an instance.
Connect to the ApsaraDB for SelectDB instance over the MySQL protocol. For more information, see Connect to an Alibaba Cloud SelectDB instance.
Create a test database and a test table.
Create a test database.
CREATE DATABASE test_db;Create a test table.
USE test_db; CREATE TABLE employees ( emp_no int NOT NULL, birth_date date, first_name varchar(20), last_name varchar(20), gender char(2), hire_date date ) UNIQUE KEY(`emp_no`) DISTRIBUTED BY HASH(`emp_no`) BUCKETS 32;
Apply for a public endpoint for the ApsaraDB for SelectDB instance. For more information, see Apply for or release a public endpoint.
Add the public IP address of your Spark environment to an IP address whitelist. For more information, see Configure an IP address whitelist.
Synchronize data from MySQL to SelectDB
Spark SQL
This example shows how to use Spark SQL to import data from an upstream MySQL database into ApsaraDB for SelectDB.
Start the
spark-sqlshell.bin/spark-sqlSubmit a job in
spark-sql.CREATE TEMPORARY VIEW mysql_tbl USING jdbc OPTIONS( "url"="jdbc:mysql://host:port/test_db", "dbtable"="employees", "driver"="com.mysql.jdbc.Driver", "user"="admin", "password"="****" ); CREATE TEMPORARY VIEW selectdb_tbl USING doris OPTIONS( "table.identifier"="test_db.employees", "fenodes"="selectdb-cn-****-public.selectdbfe.rds.aliyuncs.com:8080", "user"="admin", "password"="****", "sink.properties.format"="json" ); INSERT INTO selectdb_tbl SELECT emp_no, birth_date, first_name, last_name, gender, hire_date FROM mysql_tbl;After the Spark job completes, log on to ApsaraDB for SelectDB and verify the imported data.
DataFrame
This example shows how to use the DataFrame API to import data from an upstream MySQL database into ApsaraDB for SelectDB.
Start the
spark-shell.bin/spark-shellSubmit a job in
spark-shell.val mysqlDF = spark.read.format("jdbc") .option("url", "jdbc:mysql://host:port/test_db") .option("dbtable", "employees") .option("driver", "com.mysql.jdbc.Driver") .option("user", "admin") .option("password", "****") .load() mysqlDF.write.format("doris") .option("fenodes", "host:httpPort") .option("table.identifier", "test_db.employees") .option("user", "admin") .option("password", "****") .option("sink.batch.size", 100000) .option("sink.max-retries", 3) .option("sink.properties.format", "json") .save()After the Spark job completes, log on to ApsaraDB for SelectDB and verify the imported data.