
# 做结果表
#### 格式语法
SQL语法格式可能在不同Flink环境下有细微差异，具体以实际环境格式为准，with后面的参数名称及参数值以此文档为准。
```
create table dwsSink (
  attr_name attr_type 
  (',' attr_name attr_type)* 
  (','PRIMARY KEY (attr_name, ...) NOT ENFORCED)
)
with (
  'connector' = 'dws',
  'url' = '',
  'tableName' = '',
  'username' = '',
  'password' = ''
);
```
#### Flink SQL配置参数
Flink SQL中设置的PRIMARY KEY将自动映射到dws-client中的uniqueKeys。参数跟随client版本发布，参数功能与client一致，以下参数说明表示为最新参数。
表1数据库配置 
| 参数        | 说明                                                                  | 默认值 |
|:---|:---|:---|
| connector | flink框架区分connector参数，固定为dws。                                        | -   |
| url       | 数据库连接地址。                                                            | -   |
| username  | 配置连接用户。                                                             | -   |
| password  | 配置密码。                                                               | -   |
| tableName | 对应dws表，默认会取public模式下的表，如果非public模式则需要显式指定，格式为：**schema.tableName**。 | -   |
   
表2连接配置 
| 参数                                   | 说明                                         | 默认值           |
|:---|:---|:---|
| connectionSize                       | 初始dws-client时的并发数量。                        | 1             |
| connectionMaxUseTimeSeconds**（已废弃）** | 连接创建多少秒后强制释放（默认单位：毫秒）。                     | 3600s（一小时）    |
| connectionMaxUseTimeThreshold        | 连接创建后强制释放时间的阈值，使用固定单位秒。新参数，从2.1.0.2版本开始支持。 | 3600（一小时）     |
| connectionMaxIdleMs                  | 连接最大空闲时间，超过后将释放（单位：毫秒）。                    | 60000ms（一分钟）  |
| connectionTimeOut                    | 连接超时时间（单位：毫秒）。                             | 300000ms（五分钟） |
   
**当dws-client为2.x版本，参数全量支持在Flink SQL通过key方式配置。**下表参数为兼容1.x版本参数，当同时配置2.x和1.x参数时，2.x版本的参数值生效。
表3DWS client写入参数 
| 参数                                                                                                         | 说明                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                           | 默认值                                                                                             |
|:---|:---|:---|
| conflictStrategy   | 有主键表数据写入时主键冲突策略： - ignore：保持原数据，忽略更新数据。  - update：用新数据中非主键列更新原数据中对应列。  - replace：用新数据替换原数据。 说明： update和replace在全字段upsert时等效，在部分字段upsert时，replace相当于将数据中不包含的列设置为null。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                      | update   |
| writeMode           | 入库方式： - auto：系统自动选择。  - copy_merge：当存在主键时使用copy方式入临时表，从临时表merge至目标表；无主键时直接copy至目标表。  - copy_upsert：当存在主键时使用copy方式入临时表，从临时表upsert至目标表；无主键时直接copy至目标表。  - upsert: 有主键用upsert sql入库；无主键用insert into入库。  - UPDATE：使用update where语法更新数据，若原表无主键可选择指定uniqueKeys，指定字段不要求必须是唯一索引，但非唯一索引可能会影响性能。  - COPY_UPDATE：数据先通过copy方式入库到临时表，通过临时表加速使用update from where方式更新目标数据。  - UPDATE_AUTO：批量小于copyWriteBatchSize使用UPDATE，否则使用COPY_UPDATE。   | auto     |
| maxFlushRetryTimes                                                                                         | 在入库时最大尝试次数，次数内执行成功则不抛出异常，每次重试间隔为 1秒 \* 次数。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                   | 3                                                                                               |
| autoFlushBatchSize                                                                                         | 自动刷库的批大小（攒批大小）。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                              | 5000                                                                                            |
| autoFlushMaxInterval                                                                                       | 自动刷库的最大间隔时间（攒批时长）。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                           | 5s                                                                                              |
| copyWriteBatchSize                                                                                         | 在"writeMode == auto"下，使用copy的批大小。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                            | 1000（2.0.0.6版本及以前默认为5000）                                                                       |
| metadataCacheSeconds                                                                                       | 系统中对元数据的最大缓存时间，例如表定义信息（单位秒）。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                 | 180                                                                                             |
| copyMode                                                                                                   | copy入库格式： - CSV：将数据拼接为CSV格式入库，该方式稳定，但性能略低。  - DELIMITER：用分隔符将数据拼接，然后入库，该方式需要数据中不包含分隔符。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                               | CSV                                                                                             |
| createTempTableMode                                                                                        | 创建临时表方式： - AS  - LIKE                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                         | AS                                                                                              |
| numberAsEpochMsForDatetime                                                                                 | 如果数据库为时间类型数据源为数字类型，是否将数据当成时间戳转换为对应时间类型。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                      | false                                                                                           |
| stringToDatetimeFormat                                                                                     | 如果数据库为时间类型数据源为字符串类型，按该格式转换为时间类型，该参数配置即开启。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                    | null                                                                                            |
| updateAll                                                                                                  | upsert时set字段是否包含主键，hstore表在全列更新（set字段包含数据库所有字段）时不需要反查数据库，性能会更好。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                              | true                                                                                            |
   
表4connector参数 
| 参数                   | 说明                                                                                | 默认值                                                                                                              |
|:---|:---|:---|
| ignoreDelete         | 忽略flink任务中的delete。                                                                | false (1.0.10前默认true)   |
| ignoreNullWhenUpdate | 是否忽略flink中字段值为null的更新，只有在"conflictStrategy == update"时有效。                         | false                                                                                                            |
| sink.parallelism     | flink系统参数用于设置sink并发数量。                                                            | 跟随上游算子                                                                                                           |
| printDataPk          | 是否在connector接收到数据时打印数据主键，用于排查问题。                                                  | false                                                                                                            |
| ignoreUpdateBefore   | 忽略flink任务中的update_before，在大表局部更新时该参数一定打开，否则有update时会导致数据的其它列被设置为null，因为会先删除再写入数据。 | true                                                                                                             |
   
#### 使用Flink SQL直连DN入库
该能力依赖flink sql **DISTRIBUTE BY** 能力，mrs有提供此能力，具体请参见[Flink SQL语法增强](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_248949.html)。
connector提供UDF函数可根据分布列值计算出下游并发结合flink sql **DISTRIBUTE BY**能力实现将数据按DN分区能力，示例：
1. 需要在SQL中引入UDF。
   ```
   CREATE temporary  FUNCTION dn_hash AS 'com.huaweicloud.dws.connectors.flink.partition.DnHashFunction';
   ```
   
2. 正常写Source SQL。
   ```
   CREATE TABLE users
   (
       id         BIGINT,
       name       STRING,
       age        INT,
       text       STRING,
       created_at TIMESTAMP(3),
       updated_at TIMESTAMP(3)
   ) WITH (
         'connector' = 'datagen',
         'fields.id.kind' = 'sequence',
         'fields.id.start' = '1',
         'fields.id.end' = '1000',
         'fields.name.length' = '10',
         'fields.age.min' = '18',
         'fields.age.max' = '60',
         'fields.text.length' = '5'
         )
   ```
   
3. Sink表定义SQL中需要新增一个字段并且要求int类型值用于接收UDF计算的结果，示例中叫dn_hash。
   ```
   create table dws_users
   (
       dn_hash int,
       id         BIGINT,
       name       STRING,
       age        INT,
       text       STRING,
       created_at TIMESTAMP(3),
       updated_at TIMESTAMP(3),
       PRIMARY KEY (id) NOT ENFORCED
   ) WITH (
         'connector' = 'dws',
         'url' = '%s',
         'tableName' = 'test.users',
         'username' = '%s',
         'autoFlushBatchSize' = '50000',
         'password' = '%s'
         )
   ```
   
4. Insert into sql使用UDF获取数据下游算子信息，同时使用**DISTRIBUTE BY** 对返回结果做数据分区，数据就会按照UDF返回信息到下游指定并行度。
   ```
   insert into dws_users select 
   /*+ DISTRIBUTEBY('dn_hash') */ 
   dn_hash('test.users',10,1024, id) as dn_hash, * from users
   ```
   
 
#### UDF函数DnHashFunction参数说明
**参数格式**
dn_hash（'dws表名',sink并行度,最大并行度,DWS表中作为分布列的字段在源数据中对应的字段名称{1,}）
**参数说明**
1. 使用时上游并行度必须不多于sink并行度，DnHashFunction同样是通过进程内获取sink算子初始化的dws client实例获取到的表元数据，如果当前进程无sink算子就会导致无法获取client实例。
2. 使用后会增加一个hash算子，如果链路有多个算子处理业务，当执行hash算子后不可以再有改变数据分区的算子，否则数据会被再次分区就不能到达指定sink算子。
3. 最大并行度默认flink自动调整的，算法中需要使用，因此自动调整的无法使用，必须通过参数设置固定并把设置值作为UDF的参数，可以通过参数pipeline.max-parallelism设置或者jar方式通过API设置：
   ```
   StreamExecutionEnvironment evn = StreamExecutionEnvironment.getExecutionEnvironment();
   evn.setParallelism(1);
   evn.setMaxParallelism(1024);
   ```
   
4. 如果分布列包含多个字段，分布列的字段顺序需要保持和DWS一致，分布列支持的字段类型和dws client一致参考参数WRITE_PARTITION_POLICY，使用功能同样需要额外配置，不可自行使用。
 
#### 示例
该示例是从kafka数据源中读取数据，写入DWS结果表中，并指定攒批时间不超过10秒，每批数据最大30000条，其具体步骤如下：
1. 在DWS数据库中创建表**public.dws_orde** **r** ：
   ```
   create table public.dws_order(
     order_id VARCHAR,
     order_channel VARCHAR,
     order_time VARCHAR,
     pay_amount FLOAT8,
     real_pay FLOAT8,
     pay_time VARCHAR,
     user_id VARCHAR,
     user_name VARCHAR,
     area_id VARCHAR
     );
   ```
   
2. 消费Kafka中order_test topic中的数据作为数据源，**public.dws_order** 作为结果表，Kafka数据为JSON格式，并且字段名称和数据库字段名称一一对应：
   ```
   CREATE TABLE kafkaSource (
     order_id string,
     order_channel string,
     order_time string,
     pay_amount double,
     real_pay double,
     pay_time string,
     user_id string,
     user_name string,
     area_id string
   ) WITH (
     'connector' = 'kafka',
     'topic' = 'order_test',
     'properties.bootstrap.servers' = 'KafkaAddress1:KafkaPort,KafkaAddress2:KafkaPort',
     'properties.group.id' = 'GroupId',
     'scan.startup.mode' = 'latest-offset',
     'format' = 'json'
   );
   CREATE TABLE dwsSink (
     order_id string,
     order_channel string,
     order_time string,
     pay_amount double,
     real_pay double,
     pay_time string,
     user_id string,
     user_name string,
     area_id string
   ) WITH (
     'connector' = 'dws',
     'url' = 'jdbc:gaussdb://DWSAddress:DWSPort/DWSdbName',
     'tableName' = 'dws_order',
     'username' = 'DWSUserName',
     'password' = 'password',
     'autoFlushMaxInterval' = '10s',
     'autoFlushBatchSize' = '30000'
   );
   insert into dwsSink select * from kafkaSource;
   ```
   
3. 给Kafka写入测试数据：
   ```
   {"order_id":"202103241000000001", "order_channel":"webShop", "order_time":"2021-03-24 10:00:00", "pay_amount":"100.00", "real_pay":"100.00", "pay_time":"2021-03-24 10:02:03", "user_id":"0001", "user_name":"Alice", "area_id":"330106"}
   ```
   
4. 等10秒后在DWS表中查询结果：
   ```
   select * from dws_order
   ```
   结果如下：
   ![](https://support.huaweicloud.com/tg-dws/figure/zh-cn_image_0000001910101116.png "点击放大")
   
 
#### 使用Flink SQL直连DN基于内核动态库
1. 使用约束：
   - 该能力从**2.1.0.1** **版本**开始支持，不需要再引入UDF函数即可支持直连DN。
   
   - DN端口监听 默认情况下，管控面部署集群的时候只会让DN的端口监听在eth2网卡，CN端口会在eth1、eth2上监听。客户端可以通过eth1网卡在同一个VPC下连接或者给eth1绑定公网IP访问，eth2为内部通信网卡只用作集群节点间通信，模型如下图所示：
     ![](https://support.huaweicloud.com/tg-dws/figure/zh-cn_image_0000002454891549.png "点击放大")
     对于正常用户来说客户只能通过eth1网卡连接集群，因此要支持用户连接必须将DN的端口监听在eth1网卡，需要在每个节点上执行脚本：
     ```
     eth1IP=`ifconfig eth1|grep -w inet|awk '{print $2}'`
     path=`ps ux |grep -w gaussdb|grep -E "primary0|standby1"|awk '{print $15}'`
     for p in `echo $path` ; do
         echo $p sed -i "s/isten_addresses = '/isten_addresses = '$eth1IP,/g" /var/chroot/$p/postgresql.conf
     done
     ```
     ![](https://support.huaweicloud.com/tg-dws/public_sys-resources/note_3.0-zh-cn.png)
     - 上述配置一定要**重启**才会生效。
     
     - 如果客户端在**资源租户**下和集群在同一个VPC下则无需此设置。
       
   
   - 安全组配置 **此约束针对云上集群**。查询DWS集群对应的DN端口号（云上环境一般为40000），默认云上的集群没有放开DN端口的安全组，添加一条规则进行放通。
     ![](https://support.huaweicloud.com/tg-dws/figure/zh-cn_image_0000002562343783.png "点击放大")
     
   
   - 防火墙配置 默认防火墙不开放DN端口，放开防火墙可参考如下脚本，脚本执行需要root用户（具体脚本执行需要适配下自身环境，确保放通DN端口）：
     ```
     path=`ps -ef |grep -w gaussdb|grep -E "primary0|standby1"|awk '{print $12}'`
     for p in `echo $path` ; do
         port=`grep port /var/chroot/$p/postgresql.conf|head -n 1|awk '{print $3}'
         `sed -i "13a -A INPUT -p tcp -m tcp --dport $port -j ACCEPT" /etc/sysconfig/iptables
     done
     service iptables restart
     ```
     
   
   - 开启白名单 DWS集群DN上白名单默认只对集群内部放开，在其他节点无法访问，因此在管控面部署中需要放开该DN端口的白名单，可执行以下命令：
     ```
     gs_guc set -Z datanode -I all -N all -h "host all all 0.0.0.0/0 sha256"
     ```
     
    
2. 使用方式 在目标表的Flink SQL中添加以下两个关键的配置即可开启直连DN：
   - **dws.client.write.partition-policy** 指定为DN后，即表示开启直连DN。
   
   - **dws.client.network.internal-privateIp** 配置DN节点的IP映射，需要配置所有DN节点。internalIp为DN节点的内部IP，privateIp是客户端可以访问通的IP，格式为两个IP中间以冒号分割，不同DN之间以分号分割。
     ```
     create table dwsSink (
       attr_name attr_type 
       (',' attr_name attr_type)* 
       (','PRIMARY KEY (attr_name, ...) NOT ENFORCED)
     )
     with (
       'connector' = 'dws',
       'url' = '',
       'tableName' = '',
       'username' = '',
       'password' = '',
       'dws.client.write.partition-policy' = 'DN',
       'dws.client.network.internal-privateIp' = 'internalIp:privateIp;internalIp:privateIp;......'
     );
     ```
     
    
3. 场景限制
   - 当前直连DN特性仅支持表分布列字段为**int、bigint、varchar、text**类型，其他类型会抛出异常。
   
   - 当前直连DN特性仅支持**SQL_ASCII** 、**UTF-8**编码格式，其他类型会抛出异常。
   
   - 当前直连DN特性仅支持**Hash类型**的分布表，其他分布类型不会报错，会回退成连CN入库场景。
   
   - 当前直连DN特性不支持有**超长字符截断** 或**非法字符替换** 的配置的场景。
     具体来说，当DWS集群开启**td_compatible_truncation** ，或者配置了**load_auto_truncation** 、**change_illegal_char**后会回退成CN入库场景。
     
   
   - 当前直连DN特性**不支持** 数据同步期间DWS库执行以下操作：
     - 修改分布列。
     
     - 执行数据同步表的DDL语句。
     
     - 修改DN信息，如扩缩容等。
      
   
   - 直连DN特性对比直连CN的场景，主要适用于**CN并发度高、资源使用率较高、达到CN性能瓶颈**后的场景，此时开启直连DN性能提升较大，有明显优势。在CN未达到资源瓶颈时，由于本身MPP架构，与直连CN性能基本持平。
   
   - 直连DN使用的动态库，引入了ICU组件，需要**依赖gcc、c++组件** ，具体来说是**libstdc++** 的动态库，要求版本在**6.0.21及以上** 。如果没有，需要在环境上提前安装好，并且**CXXABI的版本需要\>=1.3.8** ，**GLIBC的版本需要\>=3.4.21**。
    
![](https://support.huaweicloud.com/tg-dws/public_sys-resources/note_3.0-zh-cn.png)
- 可以使用以下命令查看当前机器的libstdc++动态库支持什么版本的CXXABI和GLIBC（如动态库在/usr/lib64下）：
  ```
  strings /usr/lib64/libstdc++.so.6 | grep CXXABI
  strings /usr/lib64/libstdc++.so.6 | grep GLIBCXX
  ```
  
- 当遇到版本过低时，可以采用**版本升级** 或者**手动替换libstdc++动态库**的方案进行解决。
 
#### 常见问题
- Q：**writeMode** 参数设置什么值比较合适？
  A：根据业务场景分**update** （只更新存在的数据）和**upsert** （对于同一主键数据如果存在就更新，不存在就新增一条数据）两个类型，推荐直接使用**auto** 方式即可，该方式下会根据数据量的大小自动选择，如果数据量较大会增大攒批参数**autoFlushBatchSize**，即可提升入库性能。
  具体来说，当使用**auto** 作为writeMode时（默认情况下copy的数据攒批大小**copyWriteBatchSize** 为1000），低于**copyWriteBatchSize** 时会走upsert，高于**copyWriteBatchSize**时分两种情况：有主键时走copy到临时表+upsert，无主键时直接copy到目标表。
  
- Q：**autoFlushBatchSize** 和**autoFlushMaxInterval** 怎么设置比较合适？
  A：autoFlushBatchSize参数用于设置最大攒批条数，autoFlushMaxInterval参数用于设置最大攒批间隔，两个参数分别从时间和空间维度管控攒批。
  - 通过autoFlushMaxInterval可保证数据量较小时的时效性，如对时效性无强制要求通常不建议设置的太小，建议不低于3s走默认值即可。
  
  - 通过autoFlushBatchSize可控制一批数据的最大条数，一般来说攒批量越大，对于整体入库性能会更好，对性能来说通常该参数的设置推荐越大越好，参数的设置根据业务数据的大小以及flink运行内存来设置，保证不内存溢出。 对于大多业务来说无需设置autoFlushMaxInterval，将autoFlushBatchSize设置为50000即可。
    
    
- Q：遇到数据库死锁了怎么办？ A：通常出现死锁大致分为**行锁** 和**分布式死锁**。
  - 行锁：该场景通常为同一主键数据的并发更新造成行锁，该情况可以通过对数据做**key by** 解决，**key by** 必须根据数据库主键做，保证同一个主键数据会在同一个并发中，破坏掉并发更新的条件，无法造成死锁。同时需要在DWS-connector侧配置**connectionSize=1** ，保证只有一个connection进行入库，防止工具侧并发死锁。Flink SQL做**key by** 需要Flink本身支持，对于DLI/MRS均能实现，如MRS flink通过增加参数**"key-by-before-sink=true"**可实现key by。具体怎么使用以实现方为准，对于无法使用的建议使用API方式入库。
  
  - 分布式死锁：该场景通常为列存表的并发更新造成分布式死锁，暂无法解决，建议使用行存或者hstore。
   
- Q：遇到入库SQL执行超时怎么办？ A：SQL超时的报错一般包含以下内容：**canceling statement due to statement timeout**。解决这个问题可以从以下两个方面入手：
  通过DWS内核的Top SQL进行分析，查看具体的SQL和报错的SQL是否对应，这一步通常是为了确认问题。
  - 解决思路1：SQL超时可能是由DWS节点的资源不够导致，可以查看DWS资源监控详情，重点关注**CPU、IO、磁盘**等信息，若有明显波动，可能为资源不足导致。解决方案为资源升配，一次性解决问题。
  
  - 解决思路2：SQL超时可能是由于connector工具侧本身配置的默认超时时间不够执行完整条SQL（默认为5min）。此时需要将超时时间调大，直接在Flink SQL的with参数中添加以下配置即可：
  
  
  表5不同版本解决SQL超时问题参数说明 
  | DWS-connector版本 | 参数名                              | 参数值（单位：ms） | 配置示例                                      |
  |:---|:---|:---|:---|
  | 1.x             | **connectionTimeOut**            | 600000     | 'connectionTimeOut' = '600000'            |
  | 2.x             | **dws.client.timeout.statement** | 900000     | 'dws.client.timeout.statement' = '900000' |
     
  
- Q：遇到Flink的taskManager中有入库异常，但是**flink ui页面上没有展示** 怎么办？
  A：此为Flink的同步机制导致。Flink通常只感知sink算子直接抛出的异常，不会感知sink算子内部启动其他子线程抛出的异常。但是这个不影响Flink自身的重试策略。
  
- Q：**ignoreNullWhenUpdate** 这个参数什么情况下使用？
  A：**ignoreNullWhenUpdate**参数通常是为了应对部分列更新的业务而引入。基于部分列更新，入库SQL需要根据数据特征定制化SQL（即upsert时只针对有值的列进行更新），否则会存在数据覆盖的问题，导致数据不一致。但是引入该参数后可能会导致以下问题：
  - 分组过多。由于配置该参数后攒的一批数据通常存在多种数据特征，所以需要对不同的数据特征的数据进行划分，每个划分后的小组执行独立的SQL进行入库。基于此情形，可能存在极端情况，即原表为大宽表，且数据特征较散，则势必会导致分组过多，导致一批中入库SQL频繁执行，没有达到大批量执行的目的。
  
  - 部分列更新upsert性能差。此为GaussDB内核侧的性能基线---**部分列upsert要比全列upsert慢3倍以上**。
  
  
  如果遇到上述性能问题，则需要业务侧进行数据整改，或者接受不开此参数的影响。
  
- Q：涉及忽略数据前后空格的场景怎么办？ A：忽略前后空格场景可以使用Flink SQL中的**trim函数**，如 insert into sink select trim(a),trim(b) from source; 即可达到业务效果。
  
- Q：涉及源端数据（如MySQL）中的varchar字段中携带有**\\0** ，同步后发现丢失怎么办？
  A：丢失是正常现象，DWS-connector工具中默认对\\0符号进行了替换，将其替换成了空串。
  \\0一般会被认为是字符中的终止符，直接进行入库可能会造成报错，如DWS的Copy SQL默认就不支持\\0入库，会报错**invalid byte sequence for encoding "UTF8": 0x00**。
  - 规避手段：使用Upsert入库模式。但是Upsert模式一般会比较慢，需要注意Flink数据反压，甚至导致内存溢出的情况。
  
  - 解决方式：配置参数 **'dws.client.write.format.string-u0000' = 'false'** ，即可去掉上述替换空串的默认行为 。同时配置 **'** **dws.client.copy.compatible-illegal-chars.enable** **' = 'true'** ，即可解决此问题。后者参数的原理为开启此参数后，Copy语句会拼接一个**compatible_illegal_chars**参数，DWS内核会对特殊字符做容错处理，避免Copy时报错。
   
- Q：当sink业务表**无主键** ，同时指定writeMode为**copy** 时，若此时发生故障，触发Flink重试后会导致数据重复的问题，该如何解决？
  A：Flink故障后重试会基于上一次成功的Checkpoint点位进行恢复，即source算子会将消费位点回退到Checkpoint中记录的位置，那么必然会导致有一部分数据重复进行入库（**上一次Checkpoint成功时间** 到**故障发生时间**这一时间范围内）。如果此时业务表无主键，那么就会再次插入一份数据，导致数据重复，或者说导致数据不一致的产生。
  - 解决方案1：使用**Upsert** 语义。配合业务表**主键**使用，即可达到写入结果一致的效果。
  
  - 解决方案2：开启sink端的事务，部分数据源支持两阶段提交，即Flink先预导入数据，当sink算子的Checkpoint成功后再提交事务。但**此方案dws-connector暂未支持**。
   
- Q：使用COPY_UPSERT或者COPY_MERGE入库模式时，如果发现临时表频繁创建，导致系统表膨胀怎么办？ A：此问题一般是由于连接配置了最大使用时间（参数为**connectionMaxUseTimeSeconds** ），或连接最大空闲时间（参数为**connectionMaxIdleMs** ），而这两个配置的时间过短导致。**这两个参数的原理为：每次建立JDBC连接进行入库后，会检查连接的使用时间或者空闲时间是否已超过配置的阈值，如果已超过会将已有连接回收，临时表也会随之销毁**。而下一批数据入库时又会再次取新的JDBC连接，建立新的临时表，重复此步骤后就导致了问题产生。
  - 解决思路1：旧版本（2.1.0.1及以前）使用的**connectionMaxUseTimeSeconds** 和**connectionMaxIdleMs** 参数均为时间类型，不携带时间单位时均以**ms** **(毫秒)** 作为默认单位，所以请检查配置的时间是否过短或未携带相应的时间单位（**ms(毫秒)、s(秒)、m(分钟)、h(小时)、d(天)**）。
  
  - 解决思路2：新版本从2.1.0.2开始，废弃**connectionMaxUseTimeSeconds** 参数，新增**connectionMaxUseTimeThreshold** 参数，与旧参数功能一致，对使用方式进行了规范和修正。**connectionMaxUseTimeThreshold**参数表示连接创建后强制释放时间的阈值，使用固定单位秒，默认值为3600。
   
 
