Configuring a Job for Synchronizing Data from PostgreSQL to DMS for Kafka
Supported Source and Destination Database Versions
| Source Database | Destination Database |
|---|---|
| PostgreSQL database of version 9.4, 9.5, 9.6, 10, 11, 12, 13, 14, 15, or 16 | Kafka cluster (2.7 and 3.x) |
Database Account Permissions
Before you use DataArts Migration for data synchronization, ensure that the source and destination database accounts meet the requirements in the following table. The required account permissions vary depending on the synchronization task type.
| Type | Required Permissions |
|---|---|
| Source database connection account | Database CONNECT permission, schema USAGE permission, table SELECT permission, sequence SELECT permission, REPLICATION connection permission, and the permission to query pg_ls_waldir() functions. For details, see How Do I Add Additional Permissions for PostgreSQL and GaussDB Data Sources? NOTE: To add the permission to create replication connections, perform the following steps:
|
| Destination database connection account | When ciphertext access is enabled for Kafka, the account must have the permissions to publish and subscribe to topics. In other scenarios, there are no special permission requirements. NOTE: |
Supported Synchronization Objects
The following table lists the objects that can be synchronized using different links in DataArts Migration.
| Type | Note |
|---|---|
| Synchronization objects |
|
Important Notes
In addition to the constraints on supported data sources and versions, connection account permissions, and synchronization objects, you also need to pay attention to the notes in the following table.
| Type | Restriction |
|---|---|
| Database |
|
| Usage | General:
Full synchronization phase: During task startup and full data synchronization, do not perform DDL operations on the source database. Otherwise, the task may fail. Incremental synchronization phase:
Troubleshooting: If any problem occurs during task creation, startup, full synchronization, incremental synchronization, or completion, rectify the fault by referring to FAQs. |
| Other |
|
Procedure
This section uses real-time synchronization from PostgreSQL to MRS Kafka as an example to describe how to configure a real-time data migration job. Before that, ensure that you have read the instructions described in Check Before Use and completed all the preparations.
- Create a real-time migration job by following the instructions in Creating a Real-Time Migration Job and go to the job configuration page.
- Select the data connection type. Select PostgreSQL for Source and MRS Kafka for Destination. Figure 1 Selecting the data connection type
- Select a job type. The default migration type is Real-time. The migration scenario is Entire DB. Figure 2 Setting the migration job type
- Configure network resources. Select the created PostgreSQL and MRS Kafka data connections and the migration resource group for which the network connection has been configured. Figure 3 Selecting data connections and a migration resource group
If no data connection is available, click Create to go to the Manage Data Connections page of the Management Center console and click Create Data Connection to create a connection. For details, see Configuring DataArts Studio Data Connection Parameters.
If no migration resource group is available, click Create to create one. For details, see Buying a DataArts Migration Resource Group Incremental Package.
- Check the network connectivity. After the data connections and migration resource group are configured, perform the following operations to check the connectivity between the data sources and the migration resource group.
- Click Source Configuration. The system will test the connectivity of the entire migration job.
- Click Test in the source and destination and migration resource group.
If the network connectivity is abnormal, see How Do I Troubleshoot a Network Disconnection Between the Data Source and Resource Group?
- Configure source parameters.
Select the databases and tables to be synchronized based on the following table.
Table 5 Selecting the databases and tables to be synchronized Synchronization Scenario
Configuration Method
Entire DB
Select the PostgreSQL databases and tables to be migrated.Figure 4 Selecting databases and tables
Both databases and tables can be customized. You can select one database and one table, or multiple databases and tables.
- Configure destination parameters. Figure 5 Kafka destination parameters
- Destination Topic Name Rule
Configure the rule for mapping source database tables to destination Kafka topics. You can specify a fixed topic or use built-in variables to synchronize data from source tables to destination topics.
The following built-in variables are available:
- Source database name: #{source_db_name}
- Source table name: #{source_table_name}
- Kafka Partition Synchronization Policy
The following three policies are available for synchronizing source data to specified partitions of destination Kafka topics:
- To Partition 0
- To different partitions based on the hash values of database names/table names
- To different partitions based on the hash values of table primary keys
If the source has no primary key, data is synchronized to partition 0 at the destination by default.
- Partitions of New Topic
If the destination Kafka does not have the corresponding topic, the topic automatically created by DataArts Migration has three partitions.
- Destination Kafka Attributes
You can set Kafka attributes and add the properties. prefix. The job will automatically remove the prefix and transfer the attributes to the Kafka client. For details about the parameters, see the configuration descriptions in the Kafka documentation.
- Advanced settings
You can add custom attributes in the Configure Task area to enable some advanced functions. For details about the parameters, see the following table.
Figure 6 Adding custom attributes
Table 6 Advanced parameters of the job for migrating data from PostgreSQL to Kafka Parameter
Type
Default Value
Unit
Description
sink.delivery-guarantee
string
at-least-once
N/A
Semantic assurance when Flink writes data to Kafka
- at-least-once: At a checkpoint, the system waits for all data in the Kafka buffer to be confirmed by the Kafka producer. No message will be lost due to events that occur on the Kafka broker. However, duplicate messages may be generated when Flink is restarted because Flink processes old data again.
- exactly-once: In this mode, the Kafka sink writes all data through the transactions submitted at a checkpoint. Therefore, if the consumer reads only submitted data, duplicate data will not be generated when Flink is restarted. However, data is visible only when a checkpoint is complete, so you need to adjust the checkpoint interval as needed.
- Destination Topic Name Rule
- Update the mapping between the source table and destination table and check whether the mapping is correct. Figure 7 Mapping between source and destination tables
- Configure task parameters.
Table 7 Task parameters Parameter
Description
Default Value
Execution Memory
Memory allocated for job execution, which automatically changes with the number of CPU cores.
8GB
CPU Cores
Value range: 2 to 32
For each CPU core added, 4 GB execution memory and one concurrency are automatically added.
2
Concurrency
Maximum number of jobs that can be concurrently executed. This parameter does not need to be configured and automatically changes with the number of CPU cores.
1
Auto Retry
Whether to enable automatic retry upon a job failure
No
Maximum Retries
This parameter is displayed when Auto Retry is set to Yes.
1
Retry Interval (Seconds)
This parameter is displayed when Auto Retry is set to Yes.
120
Write Dirty Data
Whether to record dirty data. By default, dirty data is not recorded. If there is a large amount of dirty data, the synchronization speed of the task is affected.
- No: Dirty data is not recorded. This is the default value.
Dirty data is not allowed. If dirty data is generated during the synchronization, the task fails and exits.
- Yes: Dirty data is allowed, that is, dirty data does not affect task execution. When dirty data is allowed and its threshold is set:
- If the generated dirty data is within the threshold, the synchronization task ignores the dirty data (that is, the dirty data is not written to the destination) and is executed normally.
- If the generated dirty data exceeds the threshold, the synchronization task fails and exits. NOTE:
Criteria for determining dirty data: Dirty data is meaningless to services, is in an invalid format, or is generated when the synchronization task encounters an error. If an exception occurs when a piece of data is written to the destination, this piece of data is dirty data. Therefore, data that fails to be written is classified as dirty data.
For example, if data of the VARCHAR type at the source is written to a destination column of the INT type, dirty data cannot be written to the migration destination due to improper conversion. When configuring a synchronization task, you can configure whether to write dirty data during the synchronization and configure the number of dirty data records (maximum number of error records allowed in a single partition) to ensure task running. That is, when the number of dirty data records exceeds the threshold, the task fails and exits.
No
Dirty Data Policy
This parameter is displayed when Write Dirty Data is set to Yes. The following policies are supported:
- Do not archive: Dirty data is only recorded in job logs, but not stored.
- Archive to OBS: Dirty data is stored in OBS and printed in job logs.
Do not archive
Write Dirty Data Link
This parameter is displayed when Dirty Data Policy is set to Archive to OBS.
Only links to OBS support dirty data writes.
-
Dirty Data Directory
OBS directory to which dirty data will be written
-
Dirty Data Threshold
This parameter is only displayed when Write Dirty Data is set to Yes.
You can set the dirty data threshold as required.
NOTE:- The dirty data threshold takes effect for each concurrency. For example, if the threshold is 100 and the concurrency is 3, the maximum number of dirty data records allowed by the job is 300.
- Value -1 indicates that the number of dirty data records is not limited.
100
Enable Heartbeat Tables
Whether to enable heartbeat tables. It is disabled by default.
If heartbeat tables are enabled, the real-time migration job creates a heartbeat table at the source (if there are no such tables). Then data will be periodically updated to the table during job execution.
- Ensure that the data source instance is writable.
- Ensure that the account has the permissions to create heartbeat tables, write data to heartbeat tables, and extract data from heartbeat tables.
No
Schema/Tablespace
Schema or tablespace of the heartbeat table
test_database
Table Name
Heartbeat table name
If the table does not exist, the real-time migration job will automatically create it. Ensure that the account used for the data source connection has the permission to create tables.
You can manually create a table. For details about the table format, see Heartbeat Table Format.
test_hearbeat_table
Heartbeat Generation Interval
Interval at which the real-time migration job generates and writes heartbeat data to the heartbeat table, in seconds
10
Write Heartbeat Data to Topic
Controls if heartbeat data is written to topics. It is disabled by default.
If this option is enabled, heartbeat data will be written to a destination topic. This option takes effect only for links where Kafka is the destination.
For details about the heartbeat data format, see Heartbeat Data Formats.
No
Add Custom Attribute
You can add custom attributes to modify some job parameters and enable some advanced functions. For details, see Job Performance Optimization.
-
- No: Dirty data is not recorded. This is the default value.
- Submit and run the job.
After configuring the job, click Submit in the upper left corner to submit the job.
Figure 8 Submitting the job
After submitting the job, click Start on the job development page. In the displayed dialog box, set required parameters and click OK.
Figure 9 Starting the job
Table 8 Parameters for starting the job Parameter
Description
Synchronous Mode
- Incremental Synchronization: Incremental data synchronization starts from a specified time point.
- Full and incremental synchronization: All data is synchronized first, and then incremental data is synchronized in real time.
- Monitor the job.
On the job development page, click Monitor to go to the Job Monitoring page. You can view the status and log of the job, and configure alarm rules for the job. For details, see Real-Time Migration Job O&M.
Figure 10 Monitoring the job
Performance Optimization
If the synchronization speed is too slow, rectify the fault by referring to Job Performance Optimization.
What is your overall rating for this page?
Thank you very much for your feedback. We will continue working to improve the documentation.See the reply and handling status in My Cloud VOC.
For any further questions, feel free to contact us through the chatbot.
Chatbot