Import data by using Spark

Updated at:

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.

image

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.

    Note
    • The 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

    spark-doris-connector-2.4_2.12-1.3.2

    3.1-2.12-1.3.2

    spark-doris-connector-3.1_2.12-1.3.2

    3.2-2.12-1.3.2

    spark-doris-connector-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 jars directory of your Spark installation.

    • If you run Spark in YARN cluster mode, add the JAR package to your pre-deployment package. For example:

      1. Upload the spark-doris-connector-3.2_2.12-1.3.2.jar file 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/
      2. Add the spark-doris-connector-3.2_2.12-1.3.2.jar dependency 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: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080

table.identifier

None

Yes

The destination table in your ApsaraDB for SelectDB instance. Format: database.table. For example: test_db.test_table

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 -1 means no timeout.

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 deserialize.arrow.async is set to true.

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: csv, json, and arrow.

sink.properties.*

--

No

Parameters for the Stream Load job. For example, to specify a column separator, use 'sink.properties.column_separator' = ','. For more information about Stream Load properties, see Stream Load.

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.

sink.task.use.repartition

false

No

Specifies whether to use repartition to control the number of write partitions for SelectDB. If false (the default), coalesce is used, which can reduce parallelism if it is the only action before writing.

If true, repartition is used, which allows setting an exact number of partitions but incurs additional shuffle overhead.

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 IN expression value list for predicate pushdown. If the list count exceeds this value, Spark processes the filter instead of pushing it down.

ignore-type

None

No

A comma-separated list of field types to ignore when reading the schema from a temporary view.

Example: 'ignore-type'='bitmap,hll'

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: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080

table.identifier

None

Yes

The destination table in your ApsaraDB for SelectDB instance. Format: database.table. For example: test_db.test_table

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 -1 means no timeout.

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 deserialize.arrow.async is set to true.

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: csv, json, and arrow.

sink.properties.*

--

No

Parameters for the Stream Load job. For example, to specify a column separator, use 'sink.properties.column_separator' = ','. For more information about Stream Load properties, see Stream Load.

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.

sink.task.use.repartition

false

No

Specifies whether to use repartition to control the number of write partitions for SelectDB. If false (the default), coalesce is used, which can reduce parallelism if it is the only action before writing.

If true, repartition is used, which allows setting an exact number of partitions but incurs additional shuffle overhead.

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 IN expression value list for predicate pushdown. If the list count exceeds this value, Spark processes the filter instead of pushing it down.

ignore-type

None

No

A comma-separated list of field types to ignore when reading the schema from a temporary view.

Example: 'ignore-type'='bitmap,hll'

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.

    1. Download and decompress the Spark installation package. This example uses the spark-3.1.2-bin-hadoop3.2.tgz package.

      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.tgz
    2. Place the spark-doris-connector-3.2_2.12-1.3.2.jar file in the SPARK_HOME/jars directory.

  • Prepare the data for import. This example uses a small amount of sample data from a MySQL database.

    1. 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=utf8mb3
    2. Use DMS to generate test data. For more information, see Generate test data.

  • Configure the ApsaraDB for SelectDB instance.

    1. Create an ApsaraDB for SelectDB instance. For more information, see Create an instance.

    2. Connect to the ApsaraDB for SelectDB instance over the MySQL protocol. For more information, see Connect to an Alibaba Cloud SelectDB instance.

    3. Create a test database and a test table.

      1. Create a test database.

        CREATE DATABASE test_db;
      2. 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;
    4. Apply for a public endpoint for the ApsaraDB for SelectDB instance. For more information, see Apply for or release a public endpoint.

    5. 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.

  1. Start the spark-sql shell.

    bin/spark-sql
  2. Submit 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;
  3. 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.

  1. Start the spark-shell.

    bin/spark-shell
  2. Submit 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()
  3. After the Spark job completes, log on to ApsaraDB for SelectDB and verify the imported data.