Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Contribute to GitLab
Sign in / Register
Toggle navigation
D
dsk-dsc-flink
Project
Project
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
shezaixing
dsk-dsc-flink
Commits
773fe96f
Commit
773fe96f
authored
Aug 13, 2026
by
375138141
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
kingbase
parent
418dd22d
Changes
9
Show whitespace changes
Inline
Side-by-side
Showing
9 changed files
with
147 additions
and
125 deletions
+147
-125
pom.xml
pom.xml
+77
-98
AsyncDbDataTransferFunction.java
...link/dsc/common/function/AsyncDbDataTransferFunction.java
+4
-3
AsyncDbDataTransferFunctionNew.java
...k/dsc/common/function/AsyncDbDataTransferFunctionNew.java
+4
-3
DbDataTransferFunction.java
...dsk/flink/dsc/common/function/DbDataTransferFunction.java
+4
-3
DbDataSlideSink.java
...n/java/com/dsk/flink/dsc/common/sink/DbDataSlideSink.java
+16
-5
DbDataTransferSink.java
...ava/com/dsk/flink/dsc/common/sink/DbDataTransferSink.java
+16
-5
DbDataTransferSinkBatch.java
...om/dsk/flink/dsc/common/sink/DbDataTransferSinkBatch.java
+24
-6
EnvProperties.java
src/main/java/com/dsk/flink/dsc/utils/EnvProperties.java
+1
-1
EtlUtils.java
src/main/java/com/dsk/flink/dsc/utils/EtlUtils.java
+1
-1
No files found.
pom.xml
View file @
773fe96f
...
@@ -32,16 +32,19 @@
...
@@ -32,16 +32,19 @@
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-csv
</artifactId>
<artifactId>
flink-csv
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-json
</artifactId>
<artifactId>
flink-json
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-avro
</artifactId>
<artifactId>
flink-avro
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
io.vertx
</groupId>
<groupId>
io.vertx
</groupId>
...
@@ -71,11 +74,13 @@
...
@@ -71,11 +74,13 @@
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-table-api-java-bridge_2.12
</artifactId>
<artifactId>
flink-table-api-java-bridge_2.12
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-statebackend-rocksdb_2.12
</artifactId>
<artifactId>
flink-statebackend-rocksdb_2.12
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
com.alibaba
</groupId>
<groupId>
com.alibaba
</groupId>
...
@@ -112,21 +117,13 @@
...
@@ -112,21 +117,13 @@
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-jdbc_2.12
</artifactId>
<artifactId>
flink-jdbc_2.12
</artifactId>
<version>
${flink-jdbc-version}
</version>
<version>
${flink-jdbc-version}
</version>
<scope>
provided
</scope>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-format-common
</artifactId>
<artifactId>
flink-format-common
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
</dependency>
<scope>
provided
</scope>
<dependency>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-json
</artifactId>
<version>
${flink-version}
</version>
</dependency>
<dependency>
<groupId>
org.apache.kafka
</groupId>
<artifactId>
kafka-clients
</artifactId>
<version>
${kafka-clinet-veriosn}
</version>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.projectlombok
</groupId>
<groupId>
org.projectlombok
</groupId>
...
@@ -138,25 +135,23 @@
...
@@ -138,25 +135,23 @@
<artifactId>
hutool-all
</artifactId>
<artifactId>
hutool-all
</artifactId>
<version>
${hutool-version}
</version>
<version>
${hutool-version}
</version>
</dependency>
</dependency>
<!--<dependency>
<groupId>com.dsk.io.tidb</groupId>
<artifactId>flink-tidb-connector-1.13</artifactId>
<version>0.0.5.1</version>
</dependency>-->
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-streaming-scala_2.12
</artifactId>
<artifactId>
flink-streaming-scala_2.12
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-runtime-web_2.12
</artifactId>
<artifactId>
flink-runtime-web_2.12
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-table-api-scala-bridge_2.12
</artifactId>
<artifactId>
flink-table-api-scala-bridge_2.12
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
<exclusions>
<exclusions>
<exclusion>
<exclusion>
<artifactId>
force-shading
</artifactId>
<artifactId>
force-shading
</artifactId>
...
@@ -176,6 +171,7 @@
...
@@ -176,6 +171,7 @@
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-connector-jdbc_2.12
</artifactId>
<artifactId>
flink-connector-jdbc_2.12
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
<exclusions>
<exclusions>
<exclusion>
<exclusion>
<artifactId>
force-shading
</artifactId>
<artifactId>
force-shading
</artifactId>
...
@@ -216,12 +212,33 @@
...
@@ -216,12 +212,33 @@
<artifactId>
tikv-client-java
</artifactId>
<artifactId>
tikv-client-java
</artifactId>
<groupId>
org.tikv
</groupId>
<groupId>
org.tikv
</groupId>
</exclusion>
</exclusion>
<exclusion>
<artifactId>
connect-api
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
<exclusion>
<artifactId>
connect-runtime
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
<exclusion>
<artifactId>
connect-json
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
<exclusion>
<artifactId>
connect-transforms
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
<exclusion>
<artifactId>
kafka-tools
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
</exclusions>
</exclusions>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-java
</artifactId>
<artifactId>
flink-java
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
<exclusions>
<exclusions>
<exclusion>
<exclusion>
<artifactId>
force-shading
</artifactId>
<artifactId>
force-shading
</artifactId>
...
@@ -237,6 +254,7 @@
...
@@ -237,6 +254,7 @@
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-connector-base
</artifactId>
<artifactId>
flink-connector-base
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
<exclusions>
<exclusions>
<exclusion>
<exclusion>
<artifactId>
force-shading
</artifactId>
<artifactId>
force-shading
</artifactId>
...
@@ -244,11 +262,11 @@
...
@@ -244,11 +262,11 @@
</exclusion>
</exclusion>
</exclusions>
</exclusions>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-streaming-java_2.12
</artifactId>
<artifactId>
flink-streaming-java_2.12
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
<exclusions>
<exclusions>
<exclusion>
<exclusion>
<artifactId>
force-shading
</artifactId>
<artifactId>
force-shading
</artifactId>
...
@@ -272,6 +290,7 @@
...
@@ -272,6 +290,7 @@
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-clients_2.12
</artifactId>
<artifactId>
flink-clients_2.12
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
<exclusions>
<exclusions>
<exclusion>
<exclusion>
<artifactId>
slf4j-api
</artifactId>
<artifactId>
slf4j-api
</artifactId>
...
@@ -279,11 +298,11 @@
...
@@ -279,11 +298,11 @@
</exclusion>
</exclusion>
</exclusions>
</exclusions>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-table-common
</artifactId>
<artifactId>
flink-table-common
</artifactId>
<version>
${flink-version}
</version>
<version>
${flink-version}
</version>
<scope>
provided
</scope>
<exclusions>
<exclusions>
<exclusion>
<exclusion>
<artifactId>
force-shading
</artifactId>
<artifactId>
force-shading
</artifactId>
...
@@ -295,16 +314,11 @@
...
@@ -295,16 +314,11 @@
</exclusion>
</exclusion>
</exclusions>
</exclusions>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
com.ververica
</groupId>
<groupId>
com.ververica
</groupId>
<artifactId>
flink-connector-mysql-cdc
</artifactId>
<artifactId>
flink-connector-mysql-cdc
</artifactId>
<version>
${flink-cdc-version}
</version>
<version>
${flink-cdc-version}
</version>
<exclusions>
<exclusions>
<!--<exclusion>
<artifactId>mysql-connector-java</artifactId>
<groupId>mysql</groupId>
</exclusion>-->
<exclusion>
<exclusion>
<artifactId>
slf4j-api
</artifactId>
<artifactId>
slf4j-api
</artifactId>
<groupId>
org.slf4j
</groupId>
<groupId>
org.slf4j
</groupId>
...
@@ -313,90 +327,76 @@
...
@@ -313,90 +327,76 @@
<artifactId>
flink-connector-debezium
</artifactId>
<artifactId>
flink-connector-debezium
</artifactId>
<groupId>
com.ververica
</groupId>
<groupId>
com.ververica
</groupId>
</exclusion>
</exclusion>
<exclusion>
<artifactId>
kafka-clients
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
<exclusion>
<artifactId>
connect-api
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
<exclusion>
<artifactId>
connect-runtime
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
<exclusion>
<artifactId>
connect-json
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
<exclusion>
<artifactId>
connect-transforms
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
<exclusion>
<artifactId>
kafka-tools
</artifactId>
<groupId>
org.apache.kafka
</groupId>
</exclusion>
</exclusions>
</exclusions>
</dependency>
</dependency>
<!--<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-sql-connector-kafka_2.12</artifactId>
<version>${flink-version}</version>
</dependency>-->
<dependency>
<dependency>
<groupId>
org.apache.flink
</groupId>
<groupId>
org.apache.flink
</groupId>
<artifactId>
flink-connector-kafka_2.12
</artifactId>
<artifactId>
flink-connector-kafka_2.12
</artifactId>
<version>
${flink-version}
</version>
<version>
1.13.6
</version>
<scope>
compile
</scope>
<exclusions>
<exclusions>
<exclusion>
<exclusion>
<artifactId>
flink-connector-base
</artifactId>
<groupId>
org.apache.flink
</groupId>
</exclusion>
<exclusion>
<artifactId>
force-shading
</artifactId>
<groupId>
org.apache.flink
</groupId>
</exclusion>
<exclusion>
<artifactId>
kafka-clients
</artifactId>
<groupId>
org.apache.kafka
</groupId>
<groupId>
org.apache.kafka
</groupId>
<artifactId>
kafka-clients
</artifactId>
</exclusion>
</exclusion>
</exclusions>
</exclusions>
</dependency>
</dependency>
<dependency>
<groupId>
org.apache.kafka
</groupId>
<artifactId>
kafka-clients
</artifactId>
<version>
2.4.1
</version>
</dependency>
<dependency>
<dependency>
<groupId>
com.alibaba
</groupId>
<groupId>
com.alibaba
</groupId>
<artifactId>
fastjson
</artifactId>
<artifactId>
fastjson
</artifactId>
<version>
${fast-json-version}
</version>
<version>
${fast-json-version}
</version>
</dependency>
</dependency>
<dependency>
<dependency>
<groupId>
org.dom4j
</groupId>
<groupId>
org.dom4j
</groupId>
<artifactId>
dom4j
</artifactId>
<artifactId>
dom4j
</artifactId>
<version>
2.1.3
</version>
<version>
2.1.3
</version>
</dependency>
</dependency>
<!-- 人大金仓 Kingbase 驱动 V9R1C10 -->
<!-- 人大金仓 Kingbase 驱动 -->
<dependency>
<dependency>
<groupId>
c
om.kingbase8
</groupId>
<groupId>
c
n.com.kingbase
</groupId>
<artifactId>
kingbase8
</artifactId>
<artifactId>
kingbase8
</artifactId>
<version>
${kingbase8-version}
</version>
<version>
9.0.0
</version>
<scope>
compile
</scope>
</dependency>
</dependency>
</dependencies>
</dependencies>
<repositories>
<repository>
<id>
dskmaven
</id>
<name>
dskmaven
</name>
<url>
http://120.27.13.145:8081/nexus/content/groups/public/
</url>
</repository>
</repositories>
<build>
<build>
<finalName>
${artifactId}
</finalName>
<finalName>
${artifactId}
</finalName>
<plugins>
<plugins>
<plugin>
<groupId>
org.apache.maven.plugins
</groupId>
<artifactId>
maven-assembly-plugin
</artifactId>
<version>
3.0.0
</version>
<configuration>
<descriptorRefs>
<descriptorRef>
jar-with-dependencies
</descriptorRef>
</descriptorRefs>
</configuration>
<executions>
<execution>
<id>
make-assembly
</id>
<phase>
package
</phase>
<goals>
<goal>
single
</goal>
</goals>
</execution>
</executions>
</plugin>
<plugin>
<plugin>
<groupId>
org.apache.maven.plugins
</groupId>
<groupId>
org.apache.maven.plugins
</groupId>
<artifactId>
maven-compiler-plugin
</artifactId>
<artifactId>
maven-compiler-plugin
</artifactId>
<version>
3.10.1
</version>
<version>
3.10.1
</version>
<!-- 下面指定为自己需要的版本 -->
<configuration>
<configuration>
<encoding>
utf8
</encoding>
<encoding>
utf8
</encoding>
<source>
1.8
</source>
<source>
1.8
</source>
...
@@ -408,6 +408,7 @@
...
@@ -408,6 +408,7 @@
</excludes>
</excludes>
</configuration>
</configuration>
</plugin>
</plugin>
<!-- 只保留shade打包,删除assembly -->
<plugin>
<plugin>
<groupId>
org.apache.maven.plugins
</groupId>
<groupId>
org.apache.maven.plugins
</groupId>
<artifactId>
maven-shade-plugin
</artifactId>
<artifactId>
maven-shade-plugin
</artifactId>
...
@@ -416,52 +417,31 @@
...
@@ -416,52 +417,31 @@
<createDependencyReducedPom>
false
</createDependencyReducedPom>
<createDependencyReducedPom>
false
</createDependencyReducedPom>
</configuration>
</configuration>
<executions>
<executions>
<!-- Run shade goal on package phase -->
<execution>
<execution>
<phase>
package
</phase>
<phase>
package
</phase>
<goals>
<goals>
<goal>
shade
</goal>
<goal>
shade
</goal>
</goals>
</goals>
<configuration>
<configuration>
<!--<transformers combine.children="append">
<!– The service transformer is needed to merge META-INF/services files –>
<transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/>
<!– ... –>
</transformers>-->
<!-- 合并多个connetor 的META-INF.services 文件-->
<transformers>
<transformers>
<transformer
<transformer
implementation=
"org.apache.maven.plugins.shade.resource.AppendingTransformer"
>
implementation=
"org.apache.maven.plugins.shade.resource.AppendingTransformer"
>
<resource>
reference.conf
</resource>
<resource>
reference.conf
</resource>
</transformer>
</transformer>
<!-- The service transformer is needed to merge META-INF/services files -->
<transformer
implementation=
"org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"
/>
<transformer
<transformer
implementation=
"org.apache.maven.plugins.shade.resource.ApacheNoticeResourceTransformer"
>
implementation=
"org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"
/>
<transformer
implementation=
"org.apache.maven.plugins.shade.resource.ApacheNoticeResourceTransformer"
>
<projectName>
Apache Flink
</projectName>
<projectName>
Apache Flink
</projectName>
<encoding>
UTF-8
</encoding>
<encoding>
UTF-8
</encoding>
</transformer>
</transformer>
</transformers>
</transformers>
<!-- 自动排除不使用的类,缩小jar包体积-->
<!-- <minimizeJar>true</minimizeJar>-->
<artifactSet>
<artifactSet>
<excludes>
<excludes>
<exclude>
org.apache.flink:force-shading
</exclude>
<exclude>
org.apache.flink:force-shading
</exclude>
<exclude>
org.slf4j:*
</exclude>
<exclude>
org.slf4j:*
</exclude>
<exclude>
org.apache.logging.log4j:*
</exclude>
<exclude>
org.apache.logging.log4j:*
</exclude>
</excludes>
</excludes>
<!--<includes>
<indlude>com.aliyun.openservices:flink-log-connector</indlude>
<indlude>com.alibaba.hologres:hologres-connector-flink-1.12</indlude>
<indlude>com.google.guava:guava</indlude>
<indlude>org.projectlombok:lombok</indlude>
</includes>-->
</artifactSet>
</artifactSet>
<filters>
<filters>
<filter>
<filter>
<!-- Do not copy the signatures in the META-INF folder.
Otherwise, this might cause SecurityExceptions when using the JAR. -->
<artifact>
*:*
</artifact>
<artifact>
*:*
</artifact>
<excludes>
<excludes>
<exclude>
module-info.class
</exclude>
<exclude>
module-info.class
</exclude>
...
@@ -483,5 +463,4 @@
...
@@ -483,5 +463,4 @@
</plugin>
</plugin>
</plugins>
</plugins>
</build>
</build>
</project>
</project>
\ No newline at end of file
src/main/java/com/dsk/flink/dsc/common/function/AsyncDbDataTransferFunction.java
View file @
773fe96f
...
@@ -15,6 +15,7 @@ import org.slf4j.LoggerFactory;
...
@@ -15,6 +15,7 @@ import org.slf4j.LoggerFactory;
import
java.text.SimpleDateFormat
;
import
java.text.SimpleDateFormat
;
import
java.util.*
;
import
java.util.*
;
import
java.util.concurrent.ExecutorService
;
import
java.util.concurrent.ExecutorService
;
import
java.util.stream.Collectors
;
import
java.util.concurrent.LinkedBlockingDeque
;
import
java.util.concurrent.LinkedBlockingDeque
;
import
java.util.concurrent.ThreadPoolExecutor
;
import
java.util.concurrent.ThreadPoolExecutor
;
import
java.util.concurrent.TimeUnit
;
import
java.util.concurrent.TimeUnit
;
...
@@ -166,15 +167,15 @@ public class AsyncDbDataTransferFunction extends RichAsyncFunction<JSONObject,St
...
@@ -166,15 +167,15 @@ public class AsyncDbDataTransferFunction extends RichAsyncFunction<JSONObject,St
List
<
String
>
valueList
=
new
ArrayList
<>();
List
<
String
>
valueList
=
new
ArrayList
<>();
List
<
String
>
updateList
=
new
ArrayList
<>();
List
<
String
>
updateList
=
new
ArrayList
<>();
for
(
String
s
:
columnSet
)
{
for
(
String
s
:
columnSet
)
{
colBuilder
.
append
(
s
).
append
(
","
);
colBuilder
.
append
(
"\""
).
append
(
s
).
append
(
"\
","
);
valueList
.
add
(
getValueStringKingbase
(
dataObj
,
s
,
mysqlType
.
getString
(
s
)));
valueList
.
add
(
getValueStringKingbase
(
dataObj
,
s
,
mysqlType
.
getString
(
s
)));
updateList
.
add
(
s
+
" = EXCLUDED."
+
s
);
updateList
.
add
(
"\""
+
s
+
"\" = EXCLUDED.\""
+
s
+
"\""
);
}
}
colBuilder
.
setLength
(
colBuilder
.
length
()
-
1
);
colBuilder
.
setLength
(
colBuilder
.
length
()
-
1
);
String
valueString
=
String
.
join
(
","
,
valueList
);
String
valueString
=
String
.
join
(
","
,
valueList
);
String
updateString
=
String
.
join
(
","
,
updateList
);
String
updateString
=
String
.
join
(
","
,
updateList
);
List
<
String
>
pkList
=
new
ArrayList
<>(
pkNameSet
);
List
<
String
>
pkList
=
new
ArrayList
<>(
pkNameSet
);
String
pkString
=
String
.
join
(
","
,
pkList
);
String
pkString
=
pkList
.
stream
().
map
(
pk
->
"\""
+
pk
+
"\""
).
collect
(
Collectors
.
joining
(
","
)
);
if
(
pkList
.
isEmpty
())
{
if
(
pkList
.
isEmpty
())
{
return
String
.
format
(
"INSERT INTO %s (%s) values (%s)"
,
table
,
colBuilder
.
toString
(),
valueString
);
return
String
.
format
(
"INSERT INTO %s (%s) values (%s)"
,
table
,
colBuilder
.
toString
(),
valueString
);
}
}
...
...
src/main/java/com/dsk/flink/dsc/common/function/AsyncDbDataTransferFunctionNew.java
View file @
773fe96f
...
@@ -16,6 +16,7 @@ import org.slf4j.LoggerFactory;
...
@@ -16,6 +16,7 @@ import org.slf4j.LoggerFactory;
import
java.text.SimpleDateFormat
;
import
java.text.SimpleDateFormat
;
import
java.util.*
;
import
java.util.*
;
import
java.util.concurrent.ExecutorService
;
import
java.util.concurrent.ExecutorService
;
import
java.util.stream.Collectors
;
import
java.util.concurrent.LinkedBlockingDeque
;
import
java.util.concurrent.LinkedBlockingDeque
;
import
java.util.concurrent.ThreadPoolExecutor
;
import
java.util.concurrent.ThreadPoolExecutor
;
import
java.util.concurrent.TimeUnit
;
import
java.util.concurrent.TimeUnit
;
...
@@ -183,15 +184,15 @@ public class AsyncDbDataTransferFunctionNew extends RichAsyncFunction<JSONObject
...
@@ -183,15 +184,15 @@ public class AsyncDbDataTransferFunctionNew extends RichAsyncFunction<JSONObject
List
<
String
>
valueList
=
new
ArrayList
<>();
List
<
String
>
valueList
=
new
ArrayList
<>();
List
<
String
>
updateList
=
new
ArrayList
<>();
List
<
String
>
updateList
=
new
ArrayList
<>();
for
(
String
s
:
columnSet
)
{
for
(
String
s
:
columnSet
)
{
colBuilder
.
append
(
s
).
append
(
","
);
colBuilder
.
append
(
"\""
).
append
(
s
).
append
(
"\
","
);
valueList
.
add
(
getValueStringKingbase
(
dataObj
,
s
,
mysqlType
.
getString
(
s
)));
valueList
.
add
(
getValueStringKingbase
(
dataObj
,
s
,
mysqlType
.
getString
(
s
)));
updateList
.
add
(
s
+
" = EXCLUDED."
+
s
);
updateList
.
add
(
"\""
+
s
+
"\" = EXCLUDED.\""
+
s
+
"\""
);
}
}
colBuilder
.
setLength
(
colBuilder
.
length
()
-
1
);
colBuilder
.
setLength
(
colBuilder
.
length
()
-
1
);
String
valueString
=
String
.
join
(
","
,
valueList
);
String
valueString
=
String
.
join
(
","
,
valueList
);
String
updateString
=
String
.
join
(
","
,
updateList
);
String
updateString
=
String
.
join
(
","
,
updateList
);
List
<
String
>
pkList
=
new
ArrayList
<>(
pkNameSet
);
List
<
String
>
pkList
=
new
ArrayList
<>(
pkNameSet
);
String
pkString
=
String
.
join
(
","
,
pkList
);
String
pkString
=
pkList
.
stream
().
map
(
pk
->
"\""
+
pk
+
"\""
).
collect
(
Collectors
.
joining
(
","
)
);
if
(
pkList
.
isEmpty
())
{
if
(
pkList
.
isEmpty
())
{
return
String
.
format
(
"INSERT INTO %s (%s) values (%s)"
,
table
,
colBuilder
.
toString
(),
valueString
);
return
String
.
format
(
"INSERT INTO %s (%s) values (%s)"
,
table
,
colBuilder
.
toString
(),
valueString
);
}
}
...
...
src/main/java/com/dsk/flink/dsc/common/function/DbDataTransferFunction.java
View file @
773fe96f
...
@@ -17,6 +17,7 @@ import java.time.LocalDateTime;
...
@@ -17,6 +17,7 @@ import java.time.LocalDateTime;
import
java.time.ZoneId
;
import
java.time.ZoneId
;
import
java.time.format.DateTimeFormatter
;
import
java.time.format.DateTimeFormatter
;
import
java.util.*
;
import
java.util.*
;
import
java.util.stream.Collectors
;
/**
/**
* 重构代码
* 重构代码
* @author lww
* @author lww
...
@@ -151,16 +152,16 @@ public class DbDataTransferFunction extends ProcessFunction<JSONObject, Tuple3<S
...
@@ -151,16 +152,16 @@ public class DbDataTransferFunction extends ProcessFunction<JSONObject, Tuple3<S
List
<
String
>
valueList
=
new
ArrayList
<>();
List
<
String
>
valueList
=
new
ArrayList
<>();
List
<
String
>
updateList
=
new
ArrayList
<>();
List
<
String
>
updateList
=
new
ArrayList
<>();
for
(
String
s
:
columnSet
)
{
for
(
String
s
:
columnSet
)
{
colBuilder
.
append
(
s
).
append
(
","
);
colBuilder
.
append
(
"\""
).
append
(
s
).
append
(
"\
","
);
String
val
=
getValueStringKingbase
(
dataObj
,
s
,
mysqlType
.
getString
(
s
));
String
val
=
getValueStringKingbase
(
dataObj
,
s
,
mysqlType
.
getString
(
s
));
valueList
.
add
(
val
);
valueList
.
add
(
val
);
updateList
.
add
(
s
+
" = EXCLUDED."
+
s
);
updateList
.
add
(
"\""
+
s
+
"\" = EXCLUDED.\""
+
s
+
"\""
);
}
}
colBuilder
.
setLength
(
colBuilder
.
length
()
-
1
);
colBuilder
.
setLength
(
colBuilder
.
length
()
-
1
);
String
valueString
=
String
.
join
(
","
,
valueList
);
String
valueString
=
String
.
join
(
","
,
valueList
);
String
updateString
=
String
.
join
(
","
,
updateList
);
String
updateString
=
String
.
join
(
","
,
updateList
);
List
<
String
>
pkList
=
new
ArrayList
<>(
pkNameSet
);
List
<
String
>
pkList
=
new
ArrayList
<>(
pkNameSet
);
String
pkString
=
String
.
join
(
","
,
pkList
);
String
pkString
=
pkList
.
stream
().
map
(
pk
->
"\""
+
pk
+
"\""
).
collect
(
Collectors
.
joining
(
","
)
);
// 若没有主键,退化为普通 INSERT(重复会报错,由下游表约束保证)
// 若没有主键,退化为普通 INSERT(重复会报错,由下游表约束保证)
if
(
pkList
.
isEmpty
())
{
if
(
pkList
.
isEmpty
())
{
return
String
.
format
(
"INSERT INTO %s (%s) values (%s)"
,
table
,
colBuilder
.
toString
(),
valueString
);
return
String
.
format
(
"INSERT INTO %s (%s) values (%s)"
,
table
,
colBuilder
.
toString
(),
valueString
);
...
...
src/main/java/com/dsk/flink/dsc/common/sink/DbDataSlideSink.java
View file @
773fe96f
...
@@ -44,11 +44,11 @@ public class DbDataSlideSink extends RichSinkFunction<Tuple6<String,String,Strin
...
@@ -44,11 +44,11 @@ public class DbDataSlideSink extends RichSinkFunction<Tuple6<String,String,Strin
String
configTidbUrl
=
String
.
format
(
envProps
.
getDb_url
(),
envProps
.
getDb_host
(),
envProps
.
getDb_port
(),
envProps
.
getDb_database
());
String
configTidbUrl
=
String
.
format
(
envProps
.
getDb_url
(),
envProps
.
getDb_host
(),
envProps
.
getDb_port
(),
envProps
.
getDb_database
());
//System.out.println(configTidbUrl);
//System.out.println(configTidbUrl);
dataSource
=
new
DruidDataSource
();
dataSource
=
new
DruidDataSource
();
// 驱动类名优先取配置 db_driver,未配置则按方言推断,仍无则默认
MySQL
// 驱动类名优先取配置 db_driver,未配置则按方言推断,仍无则默认
Kingbase
String
driverClass
=
envProps
.
getDb_driver
();
String
driverClass
=
envProps
.
getDb_driver
();
if
(
StrUtil
.
isBlank
(
driverClass
))
{
if
(
StrUtil
.
isBlank
(
driverClass
))
{
driverClass
=
"
kingbase
"
.
equalsIgnoreCase
(
envProps
.
getDb_dialect
())
driverClass
=
"
mysql
"
.
equalsIgnoreCase
(
envProps
.
getDb_dialect
())
?
"com.
kingbase8.Driver"
:
"com.mysql.cj.jdbc
.Driver"
;
?
"com.
mysql.cj.jdbc.Driver"
:
"com.kingbase8
.Driver"
;
}
}
dataSource
.
setDriverClassName
(
driverClass
);
dataSource
.
setDriverClassName
(
driverClass
);
dataSource
.
setUsername
(
envProps
.
getDb_username
());
dataSource
.
setUsername
(
envProps
.
getDb_username
());
...
@@ -63,9 +63,20 @@ public class DbDataSlideSink extends RichSinkFunction<Tuple6<String,String,Strin
...
@@ -63,9 +63,20 @@ public class DbDataSlideSink extends RichSinkFunction<Tuple6<String,String,Strin
@Override
@Override
public
void
close
()
throws
Exception
{
public
void
close
()
throws
Exception
{
if
(
executorService
!=
null
)
{
executorService
.
shutdown
();
executorService
.
shutdown
();
try
{
if
(!
executorService
.
awaitTermination
(
60
,
TimeUnit
.
SECONDS
))
{
executorService
.
shutdownNow
();
}
}
catch
(
InterruptedException
e
)
{
executorService
.
shutdownNow
();
}
}
if
(
dataSource
!=
null
&&
!
dataSource
.
isClosed
())
{
dataSource
.
close
();
dataSource
.
close
();
}
}
}
@Override
@Override
public
void
invoke
(
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>
value
,
Context
context
)
throws
Exception
{
public
void
invoke
(
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>
value
,
Context
context
)
throws
Exception
{
...
...
src/main/java/com/dsk/flink/dsc/common/sink/DbDataTransferSink.java
View file @
773fe96f
...
@@ -42,11 +42,11 @@ public class DbDataTransferSink extends RichSinkFunction<String> {
...
@@ -42,11 +42,11 @@ public class DbDataTransferSink extends RichSinkFunction<String> {
String
configTidbUrl
=
String
.
format
(
envProps
.
getDb_url
(),
envProps
.
getDb_host
(),
envProps
.
getDb_port
(),
envProps
.
getDb_database
());
String
configTidbUrl
=
String
.
format
(
envProps
.
getDb_url
(),
envProps
.
getDb_host
(),
envProps
.
getDb_port
(),
envProps
.
getDb_database
());
//System.out.println(configTidbUrl);
//System.out.println(configTidbUrl);
dataSource
=
new
DruidDataSource
();
dataSource
=
new
DruidDataSource
();
// 驱动类名优先取配置 db_driver,未配置则按方言推断,仍无则默认
MySQL
// 驱动类名优先取配置 db_driver,未配置则按方言推断,仍无则默认
Kingbase
String
driverClass
=
envProps
.
getDb_driver
();
String
driverClass
=
envProps
.
getDb_driver
();
if
(
StrUtil
.
isBlank
(
driverClass
))
{
if
(
StrUtil
.
isBlank
(
driverClass
))
{
driverClass
=
"
kingbase
"
.
equalsIgnoreCase
(
envProps
.
getDb_dialect
())
driverClass
=
"
mysql
"
.
equalsIgnoreCase
(
envProps
.
getDb_dialect
())
?
"com.
kingbase8.Driver"
:
"com.mysql.cj.jdbc
.Driver"
;
?
"com.
mysql.cj.jdbc.Driver"
:
"com.kingbase8
.Driver"
;
}
}
dataSource
.
setDriverClassName
(
driverClass
);
dataSource
.
setDriverClassName
(
driverClass
);
dataSource
.
setUsername
(
envProps
.
getDb_username
());
dataSource
.
setUsername
(
envProps
.
getDb_username
());
...
@@ -61,9 +61,20 @@ public class DbDataTransferSink extends RichSinkFunction<String> {
...
@@ -61,9 +61,20 @@ public class DbDataTransferSink extends RichSinkFunction<String> {
@Override
@Override
public
void
close
()
throws
Exception
{
public
void
close
()
throws
Exception
{
if
(
executorService
!=
null
)
{
executorService
.
shutdown
();
executorService
.
shutdown
();
try
{
if
(!
executorService
.
awaitTermination
(
60
,
TimeUnit
.
SECONDS
))
{
executorService
.
shutdownNow
();
}
}
catch
(
InterruptedException
e
)
{
executorService
.
shutdownNow
();
}
}
if
(
dataSource
!=
null
&&
!
dataSource
.
isClosed
())
{
dataSource
.
close
();
dataSource
.
close
();
}
}
}
@Override
@Override
public
void
invoke
(
String
value
,
Context
context
)
throws
Exception
{
public
void
invoke
(
String
value
,
Context
context
)
throws
Exception
{
...
...
src/main/java/com/dsk/flink/dsc/common/sink/DbDataTransferSinkBatch.java
View file @
773fe96f
...
@@ -52,11 +52,11 @@ public class DbDataTransferSinkBatch extends RichSinkFunction<String> {
...
@@ -52,11 +52,11 @@ public class DbDataTransferSinkBatch extends RichSinkFunction<String> {
// 初始化获取配置
// 初始化获取配置
String
configTidbUrl
=
String
.
format
(
envProps
.
getDb_url
(),
envProps
.
getDb_host
(),
envProps
.
getDb_port
(),
envProps
.
getDb_database
());
String
configTidbUrl
=
String
.
format
(
envProps
.
getDb_url
(),
envProps
.
getDb_host
(),
envProps
.
getDb_port
(),
envProps
.
getDb_database
());
dataSource
=
new
DruidDataSource
();
dataSource
=
new
DruidDataSource
();
// 驱动类名优先取配置 db_driver,未配置则按方言推断,仍无则默认
MySQL
// 驱动类名优先取配置 db_driver,未配置则按方言推断,仍无则默认
Kingbase
String
driverClass
=
envProps
.
getDb_driver
();
String
driverClass
=
envProps
.
getDb_driver
();
if
(
StrUtil
.
isBlank
(
driverClass
))
{
if
(
StrUtil
.
isBlank
(
driverClass
))
{
driverClass
=
"
kingbase
"
.
equalsIgnoreCase
(
envProps
.
getDb_dialect
())
driverClass
=
"
mysql
"
.
equalsIgnoreCase
(
envProps
.
getDb_dialect
())
?
"com.
kingbase8.Driver"
:
"com.mysql.cj.jdbc
.Driver"
;
?
"com.
mysql.cj.jdbc.Driver"
:
"com.kingbase8
.Driver"
;
}
}
dataSource
.
setDriverClassName
(
driverClass
);
dataSource
.
setDriverClassName
(
driverClass
);
dataSource
.
setUsername
(
envProps
.
getDb_username
());
dataSource
.
setUsername
(
envProps
.
getDb_username
());
...
@@ -78,10 +78,28 @@ public class DbDataTransferSinkBatch extends RichSinkFunction<String> {
...
@@ -78,10 +78,28 @@ public class DbDataTransferSinkBatch extends RichSinkFunction<String> {
@Override
@Override
public
void
close
()
throws
Exception
{
public
void
close
()
throws
Exception
{
executorService
.
shutdown
();
if
(
scheduledExecutorService
!=
null
)
{
scheduledExecutorService
.
shutdown
();
scheduledExecutorService
.
shutdown
();
try
{
scheduledExecutorService
.
awaitTermination
(
5
,
TimeUnit
.
SECONDS
);
}
catch
(
InterruptedException
e
)
{
scheduledExecutorService
.
shutdownNow
();
}
}
if
(
executorService
!=
null
)
{
executorService
.
shutdown
();
try
{
if
(!
executorService
.
awaitTermination
(
60
,
TimeUnit
.
SECONDS
))
{
executorService
.
shutdownNow
();
}
}
catch
(
InterruptedException
e
)
{
executorService
.
shutdownNow
();
}
}
if
(
dataSource
!=
null
&&
!
dataSource
.
isClosed
())
{
dataSource
.
close
();
dataSource
.
close
();
}
}
}
@Override
@Override
public
void
invoke
(
String
value
,
Context
context
)
throws
Exception
{
public
void
invoke
(
String
value
,
Context
context
)
throws
Exception
{
...
...
src/main/java/com/dsk/flink/dsc/utils/EnvProperties.java
View file @
773fe96f
...
@@ -124,7 +124,7 @@ public class EnvProperties extends Properties {
...
@@ -124,7 +124,7 @@ public class EnvProperties extends Properties {
String
log_enable
;
String
log_enable
;
public
String
getLog_enable
()
{
public
String
getLog_enable
()
{
return
log
ical_delete
==
null
?
this
.
getProperty
(
"logical_delete"
)
:
logical_delet
e
;
return
log
_enable
==
null
?
this
.
getProperty
(
"log_enable"
)
:
log_enabl
e
;
}
}
public
void
setLog_enable
(
String
log_enable
)
{
public
void
setLog_enable
(
String
log_enable
)
{
...
...
src/main/java/com/dsk/flink/dsc/utils/EtlUtils.java
View file @
773fe96f
...
@@ -15,7 +15,7 @@ public class EtlUtils {
...
@@ -15,7 +15,7 @@ public class EtlUtils {
properties
.
setProperty
(
"auto.offset.reset"
,
"earliest"
);
properties
.
setProperty
(
"auto.offset.reset"
,
"earliest"
);
properties
.
setProperty
(
"sasl.jaas.config"
,
getSaslJaasConfig
(
username
,
password
));
properties
.
setProperty
(
"sasl.jaas.config"
,
getSaslJaasConfig
(
username
,
password
));
properties
.
setProperty
(
"security.protocol"
,
"SASL_PLAINTEXT"
);
properties
.
setProperty
(
"security.protocol"
,
"SASL_PLAINTEXT"
);
properties
.
setProperty
(
"sasl.mechanism"
,
"SCRAM-SHA-
512
"
);
properties
.
setProperty
(
"sasl.mechanism"
,
"SCRAM-SHA-
256
"
);
properties
.
setProperty
(
"fetch.max.bytes"
,
"20971520"
);
//20M
properties
.
setProperty
(
"fetch.max.bytes"
,
"20971520"
);
//20M
properties
.
setProperty
(
"flink.consumer.max.fetch.size"
,
"20971520"
);
//20M
properties
.
setProperty
(
"flink.consumer.max.fetch.size"
,
"20971520"
);
//20M
properties
.
setProperty
(
"session.timeout.ms"
,
"60000"
);
properties
.
setProperty
(
"session.timeout.ms"
,
"60000"
);
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment