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
418dd22d
Commit
418dd22d
authored
Aug 04, 2026
by
375138141
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
kingbase
parent
13c5b985
Changes
9
Hide whitespace changes
Inline
Side-by-side
Showing
9 changed files
with
263 additions
and
46 deletions
+263
-46
pom.xml
pom.xml
+8
-0
AsyncDbDataTransferFunction.java
...link/dsc/common/function/AsyncDbDataTransferFunction.java
+54
-6
AsyncDbDataTransferFunctionNew.java
...k/dsc/common/function/AsyncDbDataTransferFunctionNew.java
+63
-9
DbDataTransferFunction.java
...dsk/flink/dsc/common/function/DbDataTransferFunction.java
+70
-9
DbDataSlideSink.java
...n/java/com/dsk/flink/dsc/common/sink/DbDataSlideSink.java
+13
-6
DbDataTransferSink.java
...ava/com/dsk/flink/dsc/common/sink/DbDataTransferSink.java
+12
-5
DbDataTransferSinkBatch.java
...om/dsk/flink/dsc/common/sink/DbDataTransferSinkBatch.java
+15
-5
SyncCustomerDataSource.java
...n/java/com/dsk/flink/dsc/sync/SyncCustomerDataSource.java
+6
-6
EnvProperties.java
src/main/java/com/dsk/flink/dsc/utils/EnvProperties.java
+22
-0
No files found.
pom.xml
View file @
418dd22d
...
@@ -16,6 +16,7 @@
...
@@ -16,6 +16,7 @@
<fast-json-version>
1.2.75
</fast-json-version>
<fast-json-version>
1.2.75
</fast-json-version>
<kafka-clinet-veriosn>
2.7.0
</kafka-clinet-veriosn>
<kafka-clinet-veriosn>
2.7.0
</kafka-clinet-veriosn>
<mysql-connector-version>
8.0.17
</mysql-connector-version>
<mysql-connector-version>
8.0.17
</mysql-connector-version>
<kingbase8-version>
8.6.0
</kingbase8-version>
<flink-jdbc-version>
1.10.0
</flink-jdbc-version>
<flink-jdbc-version>
1.10.0
</flink-jdbc-version>
<redisson-version>
3.17.7
</redisson-version>
<redisson-version>
3.17.7
</redisson-version>
<project.build.sourceEncoding>
UTF8
</project.build.sourceEncoding>
<project.build.sourceEncoding>
UTF8
</project.build.sourceEncoding>
...
@@ -351,6 +352,13 @@
...
@@ -351,6 +352,13 @@
<version>
2.1.3
</version>
<version>
2.1.3
</version>
</dependency>
</dependency>
<!-- 人大金仓 Kingbase 驱动 -->
<dependency>
<groupId>
com.kingbase8
</groupId>
<artifactId>
kingbase8
</artifactId>
<version>
${kingbase8-version}
</version>
</dependency>
</dependencies>
</dependencies>
<repositories>
<repositories>
...
...
src/main/java/com/dsk/flink/dsc/common/function/Async
Mysql
DataTransferFunction.java
→
src/main/java/com/dsk/flink/dsc/common/function/Async
Db
DataTransferFunction.java
View file @
418dd22d
...
@@ -25,16 +25,16 @@ import java.util.concurrent.TimeUnit;
...
@@ -25,16 +25,16 @@ import java.util.concurrent.TimeUnit;
* @description mysql/tidb sink 数据组装 sql
* @description mysql/tidb sink 数据组装 sql
*
*
*/
*/
public
class
Async
Mysql
DataTransferFunction
extends
RichAsyncFunction
<
JSONObject
,
String
>
{
public
class
Async
Db
DataTransferFunction
extends
RichAsyncFunction
<
JSONObject
,
String
>
{
static
Logger
logger
=
LoggerFactory
.
getLogger
(
Async
Mysql
DataTransferFunction
.
class
);
static
Logger
logger
=
LoggerFactory
.
getLogger
(
Async
Db
DataTransferFunction
.
class
);
//数据库连接信息
//数据库连接信息
EnvProperties
dbInfoMap
;
EnvProperties
dbInfoMap
;
//线程池
//线程池
private
transient
ExecutorService
executorService
;
private
transient
ExecutorService
executorService
;
public
Async
Mysql
DataTransferFunction
(
EnvProperties
dbInfoMap
)
{
public
Async
Db
DataTransferFunction
(
EnvProperties
dbInfoMap
)
{
this
.
dbInfoMap
=
dbInfoMap
;
this
.
dbInfoMap
=
dbInfoMap
;
}
}
...
@@ -87,17 +87,23 @@ public class AsyncMysqlDataTransferFunction extends RichAsyncFunction<JSONObject
...
@@ -87,17 +87,23 @@ public class AsyncMysqlDataTransferFunction extends RichAsyncFunction<JSONObject
dataObj
.
put
(
"is_del"
,
"DELETE"
.
equals
(
type
)
?
1
:
0
);
dataObj
.
put
(
"is_del"
,
"DELETE"
.
equals
(
type
)
?
1
:
0
);
}
}
String
dialect
=
MapUtil
.
getStr
(
dbInfoMap
,
"db_dialect"
,
"mysql"
);
if
(
"INSERT"
.
equals
(
type
)){
if
(
"INSERT"
.
equals
(
type
)){
excueteSql
=
tranferInsertSql
(
table
,
dataObj
,
mysqlType
);
excueteSql
=
"kingbase"
.
equalsIgnoreCase
(
dialect
)
?
transferUpsertSqlKingbase
(
table
,
dataObj
,
mysqlType
,
pkNameSet
)
:
tranferInsertSql
(
table
,
dataObj
,
mysqlType
);
}
}
if
(
"UPDATE"
.
equals
(
type
)){
if
(
"UPDATE"
.
equals
(
type
)){
JSONObject
oldDataObj
=
oldDataList
.
getJSONObject
(
0
);
JSONObject
oldDataObj
=
oldDataList
.
getJSONObject
(
0
);
// excueteSql = tranferUpdateSql(table,dataObj,oldDataObj,mysqlType,pkNameSet);
// excueteSql = tranferUpdateSql(table,dataObj,oldDataObj,mysqlType,pkNameSet);
excueteSql
=
tranferInsertSql
(
table
,
dataObj
,
mysqlType
);
excueteSql
=
"kingbase"
.
equalsIgnoreCase
(
dialect
)
?
transferUpsertSqlKingbase
(
table
,
dataObj
,
mysqlType
,
pkNameSet
)
:
tranferInsertSql
(
table
,
dataObj
,
mysqlType
);
}
}
if
(
"DELETE"
.
equals
(
type
)){
if
(
"DELETE"
.
equals
(
type
)){
excueteSql
=
logicalDelete
?
tranferInsertSql
(
table
,
dataObj
,
mysqlType
)
:
transferDeleteSql
(
table
,
dataObj
,
mysqlType
,
pkNameSet
);
excueteSql
=
logicalDelete
?
(
"kingbase"
.
equalsIgnoreCase
(
dialect
)
?
transferUpsertSqlKingbase
(
table
,
dataObj
,
mysqlType
,
pkNameSet
)
:
tranferInsertSql
(
table
,
dataObj
,
mysqlType
))
:
transferDeleteSql
(
table
,
dataObj
,
mysqlType
,
pkNameSet
);
}
}
//处理先后顺序
//处理先后顺序
...
@@ -151,6 +157,31 @@ public class AsyncMysqlDataTransferFunction extends RichAsyncFunction<JSONObject
...
@@ -151,6 +157,31 @@ public class AsyncMysqlDataTransferFunction extends RichAsyncFunction<JSONObject
return
String
.
format
(
"INSERT INTO %s (%s) values (%s) ON DUPLICATE KEY UPDATE %s;"
,
table
,
columnString
,
valueString
,
updateString
);
return
String
.
format
(
"INSERT INTO %s (%s) values (%s) ON DUPLICATE KEY UPDATE %s;"
,
table
,
columnString
,
valueString
,
updateString
);
}
}
/**
* Kingbase(PostgreSQL 系)写法:INSERT ... ON CONFLICT (pk) DO UPDATE
*/
private
static
String
transferUpsertSqlKingbase
(
String
table
,
JSONObject
dataObj
,
JSONObject
mysqlType
,
Set
<
String
>
pkNameSet
)
{
Set
<
String
>
columnSet
=
mysqlType
.
keySet
();
StringBuilder
colBuilder
=
new
StringBuilder
();
List
<
String
>
valueList
=
new
ArrayList
<>();
List
<
String
>
updateList
=
new
ArrayList
<>();
for
(
String
s
:
columnSet
)
{
colBuilder
.
append
(
s
).
append
(
","
);
valueList
.
add
(
getValueStringKingbase
(
dataObj
,
s
,
mysqlType
.
getString
(
s
)));
updateList
.
add
(
s
+
" = EXCLUDED."
+
s
);
}
colBuilder
.
setLength
(
colBuilder
.
length
()
-
1
);
String
valueString
=
String
.
join
(
","
,
valueList
);
String
updateString
=
String
.
join
(
","
,
updateList
);
List
<
String
>
pkList
=
new
ArrayList
<>(
pkNameSet
);
String
pkString
=
String
.
join
(
","
,
pkList
);
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) ON CONFLICT (%s) DO UPDATE SET %s"
,
table
,
colBuilder
.
toString
(),
valueString
,
pkString
,
updateString
);
}
private
String
tranferUpdateSql
(
String
table
,
JSONObject
dataObj
,
JSONObject
oldDataObj
,
JSONObject
mysqlType
,
Set
<
String
>
pkNameSet
)
{
private
String
tranferUpdateSql
(
String
table
,
JSONObject
dataObj
,
JSONObject
oldDataObj
,
JSONObject
mysqlType
,
Set
<
String
>
pkNameSet
)
{
Set
<
String
>
columnSet
=
mysqlType
.
keySet
();
Set
<
String
>
columnSet
=
mysqlType
.
keySet
();
...
@@ -208,6 +239,23 @@ public class AsyncMysqlDataTransferFunction extends RichAsyncFunction<JSONObject
...
@@ -208,6 +239,23 @@ public class AsyncMysqlDataTransferFunction extends RichAsyncFunction<JSONObject
return
dataObj
.
getString
(
columnKey
);
return
dataObj
.
getString
(
columnKey
);
}
}
/**
* Kingbase(PostgreSQL 系) 取值:字符串加单引号,内部单引号转义为 ''(而非 \')
*/
private
static
String
getValueStringKingbase
(
JSONObject
dataObj
,
String
columnKey
,
String
mysqlType
)
{
if
(
null
==
dataObj
.
get
(
columnKey
))
{
return
"null"
;
}
if
(
Arrays
.
asList
(
STR_SQL_TYPE
).
contains
(
mysqlType
.
toUpperCase
()))
{
return
"'"
+
dataObj
.
getString
(
columnKey
).
replace
(
"'"
,
"''"
)
+
"'"
;
}
if
(
"DATE"
.
equalsIgnoreCase
(
mysqlType
)
||
"DATETIME"
.
equalsIgnoreCase
(
mysqlType
))
{
SimpleDateFormat
df
=
"DATETIME"
.
equalsIgnoreCase
(
mysqlType
)
?
new
SimpleDateFormat
(
"yyyy-MM-dd HH:mm:ss"
)
:
new
SimpleDateFormat
(
"yyyy-MM-dd"
);
return
String
.
format
(
"'%s'"
,
df
.
format
(
dataObj
.
getDate
(
columnKey
)));
}
return
dataObj
.
getString
(
columnKey
);
}
public
static
void
main
(
String
[]
args
)
{
public
static
void
main
(
String
[]
args
)
{
JSONObject
jsonObject
=
new
JSONObject
();
JSONObject
jsonObject
=
new
JSONObject
();
...
...
src/main/java/com/dsk/flink/dsc/common/function/Async
Mysql
DataTransferFunctionNew.java
→
src/main/java/com/dsk/flink/dsc/common/function/Async
Db
DataTransferFunctionNew.java
View file @
418dd22d
...
@@ -20,16 +20,16 @@ import java.util.concurrent.LinkedBlockingDeque;
...
@@ -20,16 +20,16 @@ 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
;
public
class
Async
Mysql
DataTransferFunctionNew
extends
RichAsyncFunction
<
JSONObject
,
Tuple3
<
String
,
String
,
Long
>>
{
public
class
Async
Db
DataTransferFunctionNew
extends
RichAsyncFunction
<
JSONObject
,
Tuple3
<
String
,
String
,
Long
>>
{
//static Logger logger = LoggerFactory.getLogger(Async
Mysql
DataTransferFunctionNew.class);
//static Logger logger = LoggerFactory.getLogger(Async
Db
DataTransferFunctionNew.class);
//数据库连接信息
//数据库连接信息
EnvProperties
dbInfoMap
;
EnvProperties
dbInfoMap
;
//线程池
//线程池
private
transient
ExecutorService
executorService
;
private
transient
ExecutorService
executorService
;
public
Async
Mysql
DataTransferFunctionNew
(
EnvProperties
dbInfoMap
)
{
public
Async
Db
DataTransferFunctionNew
(
EnvProperties
dbInfoMap
)
{
this
.
dbInfoMap
=
dbInfoMap
;
this
.
dbInfoMap
=
dbInfoMap
;
}
}
...
@@ -94,17 +94,23 @@ public class AsyncMysqlDataTransferFunctionNew extends RichAsyncFunction<JSONObj
...
@@ -94,17 +94,23 @@ public class AsyncMysqlDataTransferFunctionNew extends RichAsyncFunction<JSONObj
groupKey
=
table
.
concat
(
"-"
).
concat
(
pkValue
);
groupKey
=
table
.
concat
(
"-"
).
concat
(
pkValue
);
}
}
String
dialect
=
MapUtil
.
getStr
(
dbInfoMap
,
"db_dialect"
,
"mysql"
);
if
(
"INSERT"
.
equals
(
type
)){
if
(
"INSERT"
.
equals
(
type
)){
excueteSql
=
tranferInsertSql
(
table
,
dataObj
,
mysqlType
);
excueteSql
=
"kingbase"
.
equalsIgnoreCase
(
dialect
)
?
transferUpsertSqlKingbase
(
table
,
dataObj
,
mysqlType
,
pkNameSet
)
:
tranferInsertSql
(
table
,
dataObj
,
mysqlType
);
}
}
if
(
"UPDATE"
.
equals
(
type
)){
if
(
"UPDATE"
.
equals
(
type
)){
//JSONObject oldDataObj = oldDataList.getJSONObject(0);
//JSONObject oldDataObj = oldDataList.getJSONObject(0);
//excueteSql = tranferUpdateSql(table,dataObj,oldDataObj,mysqlType,pkNameSet);
//excueteSql = tranferUpdateSql(table,dataObj,oldDataObj,mysqlType,pkNameSet);
excueteSql
=
tranferInsertSql
(
table
,
dataObj
,
mysqlType
);
excueteSql
=
"kingbase"
.
equalsIgnoreCase
(
dialect
)
?
transferUpsertSqlKingbase
(
table
,
dataObj
,
mysqlType
,
pkNameSet
)
:
tranferInsertSql
(
table
,
dataObj
,
mysqlType
);
}
}
if
(
"DELETE"
.
equals
(
type
)){
if
(
"DELETE"
.
equals
(
type
)){
excueteSql
=
logicalDelete
?
tranferInsertSql
(
table
,
dataObj
,
mysqlType
)
:
transferDeleteSql
(
table
,
dataObj
,
mysqlType
,
pkNameSet
);
excueteSql
=
logicalDelete
?
(
"kingbase"
.
equalsIgnoreCase
(
dialect
)
?
transferUpsertSqlKingbase
(
table
,
dataObj
,
mysqlType
,
pkNameSet
)
:
tranferInsertSql
(
table
,
dataObj
,
mysqlType
))
:
transferDeleteSql
(
table
,
dataObj
,
mysqlType
,
pkNameSet
);
}
}
resultList
.
add
(
Tuple3
.
of
(
excueteSql
,
groupKey
,
ts
));
resultList
.
add
(
Tuple3
.
of
(
excueteSql
,
groupKey
,
ts
));
Boolean
logEnable
=
MapUtil
.
getBool
(
dbInfoMap
,
"log_enable"
,
false
);
Boolean
logEnable
=
MapUtil
.
getBool
(
dbInfoMap
,
"log_enable"
,
false
);
...
@@ -122,7 +128,8 @@ public class AsyncMysqlDataTransferFunctionNew extends RichAsyncFunction<JSONObj
...
@@ -122,7 +128,8 @@ public class AsyncMysqlDataTransferFunctionNew extends RichAsyncFunction<JSONObj
});
});
}
}
private
static
String
logSqlFormat
=
"INSERT INTO dsc_cdc_log (`table`,op_type,pk_columns,pk_values,data_json,cdc_ts) values ('%s','%s','%s','%s','%s', %d)"
;
private
static
String
logSqlFormatMysql
=
"INSERT INTO dsc_cdc_log (`table`,op_type,pk_columns,pk_values,data_json,cdc_ts) values ('%s','%s','%s','%s','%s', %d)"
;
private
static
String
logSqlFormatKingbase
=
"INSERT INTO dsc_cdc_log (table,op_type,pk_columns,pk_values,data_json,cdc_ts) values ('%s','%s','%s','%s','%s', %d)"
;
private
String
buildLogData
(
String
type
,
String
table
,
Set
<
String
>
pkNameSet
,
JSONObject
dataObj
,
long
ts
,
String
dataJsonStr
)
{
private
String
buildLogData
(
String
type
,
String
table
,
Set
<
String
>
pkNameSet
,
JSONObject
dataObj
,
long
ts
,
String
dataJsonStr
)
{
List
<
String
>
pkValueList
=
new
ArrayList
<>();
List
<
String
>
pkValueList
=
new
ArrayList
<>();
...
@@ -131,8 +138,13 @@ public class AsyncMysqlDataTransferFunctionNew extends RichAsyncFunction<JSONObj
...
@@ -131,8 +138,13 @@ public class AsyncMysqlDataTransferFunctionNew extends RichAsyncFunction<JSONObj
}
}
String
pkColumns
=
String
.
join
(
","
,
pkNameSet
);
String
pkColumns
=
String
.
join
(
","
,
pkNameSet
);
String
pkValues
=
String
.
join
(
"-"
,
pkValueList
);
String
pkValues
=
String
.
join
(
"-"
,
pkValueList
);
dataJsonStr
=
dataJsonStr
.
replace
(
"\\"
,
"\\\\"
);
boolean
kingbase
=
"kingbase"
.
equalsIgnoreCase
(
MapUtil
.
getStr
(
dbInfoMap
,
"db_dialect"
,
"mysql"
));
return
String
.
format
(
logSqlFormat
,
table
,
type
,
pkColumns
,
pkValues
,
dataJsonStr
,
ts
);
// Kingbase 单引号转义用 '',MySQL 用 \'
dataJsonStr
=
kingbase
?
dataJsonStr
.
replace
(
"'"
,
"''"
)
:
dataJsonStr
.
replace
(
"\\"
,
"\\\\"
).
replace
(
"'"
,
"\\'"
);
String
format
=
kingbase
?
logSqlFormatKingbase
:
logSqlFormatMysql
;
return
String
.
format
(
format
,
table
,
type
,
pkColumns
,
pkValues
,
dataJsonStr
,
ts
);
}
}
...
@@ -162,6 +174,31 @@ public class AsyncMysqlDataTransferFunctionNew extends RichAsyncFunction<JSONObj
...
@@ -162,6 +174,31 @@ public class AsyncMysqlDataTransferFunctionNew extends RichAsyncFunction<JSONObj
return
String
.
format
(
"REPLACE INTO %s (%s) values (%s);"
,
table
,
columnString
,
valueString
);
return
String
.
format
(
"REPLACE INTO %s (%s) values (%s);"
,
table
,
columnString
,
valueString
);
}
}
/**
* Kingbase(PostgreSQL 系)写法:INSERT ... ON CONFLICT (pk) DO UPDATE
*/
private
static
String
transferUpsertSqlKingbase
(
String
table
,
JSONObject
dataObj
,
JSONObject
mysqlType
,
Set
<
String
>
pkNameSet
)
{
Set
<
String
>
columnSet
=
mysqlType
.
keySet
();
StringBuilder
colBuilder
=
new
StringBuilder
();
List
<
String
>
valueList
=
new
ArrayList
<>();
List
<
String
>
updateList
=
new
ArrayList
<>();
for
(
String
s
:
columnSet
)
{
colBuilder
.
append
(
s
).
append
(
","
);
valueList
.
add
(
getValueStringKingbase
(
dataObj
,
s
,
mysqlType
.
getString
(
s
)));
updateList
.
add
(
s
+
" = EXCLUDED."
+
s
);
}
colBuilder
.
setLength
(
colBuilder
.
length
()
-
1
);
String
valueString
=
String
.
join
(
","
,
valueList
);
String
updateString
=
String
.
join
(
","
,
updateList
);
List
<
String
>
pkList
=
new
ArrayList
<>(
pkNameSet
);
String
pkString
=
String
.
join
(
","
,
pkList
);
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) ON CONFLICT (%s) DO UPDATE SET %s"
,
table
,
colBuilder
.
toString
(),
valueString
,
pkString
,
updateString
);
}
private
String
tranferUpdateSql
(
String
table
,
JSONObject
dataObj
,
JSONObject
oldDataObj
,
JSONObject
mysqlType
,
Set
<
String
>
pkNameSet
)
{
private
String
tranferUpdateSql
(
String
table
,
JSONObject
dataObj
,
JSONObject
oldDataObj
,
JSONObject
mysqlType
,
Set
<
String
>
pkNameSet
)
{
Set
<
String
>
columnSet
=
mysqlType
.
keySet
();
Set
<
String
>
columnSet
=
mysqlType
.
keySet
();
...
@@ -219,4 +256,21 @@ public class AsyncMysqlDataTransferFunctionNew extends RichAsyncFunction<JSONObj
...
@@ -219,4 +256,21 @@ public class AsyncMysqlDataTransferFunctionNew extends RichAsyncFunction<JSONObj
return
dataObj
.
getString
(
columnKey
);
return
dataObj
.
getString
(
columnKey
);
}
}
/**
* Kingbase(PostgreSQL 系) 取值:字符串加单引号,内部单引号转义为 ''(而非 \')
*/
private
static
String
getValueStringKingbase
(
JSONObject
dataObj
,
String
columnKey
,
String
mysqlType
)
{
if
(
null
==
dataObj
.
get
(
columnKey
))
{
return
"null"
;
}
if
(
Arrays
.
asList
(
STR_SQL_TYPE
).
contains
(
mysqlType
.
toUpperCase
()))
{
return
"'"
+
dataObj
.
getString
(
columnKey
).
replace
(
"'"
,
"''"
)
+
"'"
;
}
if
(
"DATE"
.
equalsIgnoreCase
(
mysqlType
)
||
"DATETIME"
.
equalsIgnoreCase
(
mysqlType
))
{
SimpleDateFormat
df
=
"DATETIME"
.
equalsIgnoreCase
(
mysqlType
)
?
new
SimpleDateFormat
(
"yyyy-MM-dd HH:mm:ss"
)
:
new
SimpleDateFormat
(
"yyyy-MM-dd"
);
return
String
.
format
(
"'%s'"
,
df
.
format
(
dataObj
.
getDate
(
columnKey
)));
}
return
dataObj
.
getString
(
columnKey
);
}
}
}
src/main/java/com/dsk/flink/dsc/common/function/
Mysql
DataTransferFunction.java
→
src/main/java/com/dsk/flink/dsc/common/function/
Db
DataTransferFunction.java
View file @
418dd22d
...
@@ -22,7 +22,7 @@ import java.util.*;
...
@@ -22,7 +22,7 @@ import java.util.*;
* @author lww
* @author lww
* @date 2025-01-14
* @date 2025-01-14
*/
*/
public
class
Mysql
DataTransferFunction
extends
ProcessFunction
<
JSONObject
,
Tuple3
<
String
,
String
,
Long
>>
{
public
class
Db
DataTransferFunction
extends
ProcessFunction
<
JSONObject
,
Tuple3
<
String
,
String
,
Long
>>
{
private
static
final
Map
<
String
,
Integer
>
STR_SQL_TYPE
;
private
static
final
Map
<
String
,
Integer
>
STR_SQL_TYPE
;
private
final
EnvProperties
dbInfoMap
;
private
final
EnvProperties
dbInfoMap
;
...
@@ -45,7 +45,7 @@ public class MysqlDataTransferFunction extends ProcessFunction<JSONObject, Tuple
...
@@ -45,7 +45,7 @@ public class MysqlDataTransferFunction extends ProcessFunction<JSONObject, Tuple
STR_SQL_TYPE
.
put
(
"JSON"
,
1
);
STR_SQL_TYPE
.
put
(
"JSON"
,
1
);
}
}
public
Mysql
DataTransferFunction
(
EnvProperties
envProps
,
OutputTag
<
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>>
logSlideTag
)
{
public
Db
DataTransferFunction
(
EnvProperties
envProps
,
OutputTag
<
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>>
logSlideTag
)
{
this
.
dbInfoMap
=
envProps
;
this
.
dbInfoMap
=
envProps
;
this
.
logSlideTag
=
logSlideTag
;
this
.
logSlideTag
=
logSlideTag
;
}
}
...
@@ -95,25 +95,35 @@ public class MysqlDataTransferFunction extends ProcessFunction<JSONObject, Tuple
...
@@ -95,25 +95,35 @@ public class MysqlDataTransferFunction extends ProcessFunction<JSONObject, Tuple
table
=
table
.
concat
(
"_"
).
concat
(
String
.
valueOf
(
val
%
i
));
table
=
table
.
concat
(
"_"
).
concat
(
String
.
valueOf
(
val
%
i
));
}
}
}
}
String
dialect
=
MapUtil
.
getStr
(
dbInfoMap
,
"db_dialect"
,
"mysql"
);
if
(
"INSERT"
.
equals
(
type
)
||
"UPDATE"
.
equals
(
type
)){
if
(
"INSERT"
.
equals
(
type
)
||
"UPDATE"
.
equals
(
type
)){
excueteSql
=
tranferInsertSql
(
table
,
dataObj
,
mysqlType
);
excueteSql
=
"kingbase"
.
equalsIgnoreCase
(
dialect
)
?
transferUpsertSqlKingbase
(
table
,
dataObj
,
mysqlType
,
pkNameSet
)
:
tranferInsertSql
(
table
,
dataObj
,
mysqlType
);
}
else
{
}
else
{
excueteSql
=
logicalDelete
?
tranferInsertSql
(
table
,
dataObj
,
mysqlType
)
:
transferDeleteSql
(
table
,
dataObj
,
mysqlType
,
pkNameSet
);
excueteSql
=
logicalDelete
?
(
"kingbase"
.
equalsIgnoreCase
(
dialect
)
?
transferUpsertSqlKingbase
(
table
,
dataObj
,
mysqlType
,
pkNameSet
)
:
tranferInsertSql
(
table
,
dataObj
,
mysqlType
))
:
transferDeleteSql
(
table
,
dataObj
,
mysqlType
,
pkNameSet
);
}
}
if
(
MapUtil
.
getBool
(
dbInfoMap
,
"log_enable"
,
false
)){
if
(
MapUtil
.
getBool
(
dbInfoMap
,
"log_enable"
,
false
)){
ctx
.
output
(
logSlideTag
,
buildLogData
(
type
,
table
,
pkColumns
,
pkColumnVals
,
ts
,
value
.
toJSONString
()));
ctx
.
output
(
logSlideTag
,
buildLogData
(
dbInfoMap
,
type
,
table
,
pkColumns
,
pkColumnVals
,
ts
,
value
.
toJSONString
()));
}
}
out
.
collect
(
Tuple3
.
of
(
excueteSql
,
groupKey
,
ts
));
out
.
collect
(
Tuple3
.
of
(
excueteSql
,
groupKey
,
ts
));
}
}
private
static
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>
buildLogData
(
String
type
,
String
table
,
StringBuilder
pkColumns
,
StringBuilder
pkValues
,
long
ts
,
String
dataJsonStr
)
{
private
static
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>
buildLogData
(
EnvProperties
dbInfoMap
,
String
type
,
String
table
,
StringBuilder
pkColumns
,
StringBuilder
pkValues
,
long
ts
,
String
dataJsonStr
)
{
if
(
pkColumns
.
length
()
>
0
)
{
if
(
pkColumns
.
length
()
>
0
)
{
pkColumns
.
setLength
(
pkColumns
.
length
()-
1
);
pkColumns
.
setLength
(
pkColumns
.
length
()-
1
);
pkValues
.
setLength
(
pkValues
.
length
()-
1
);
pkValues
.
setLength
(
pkValues
.
length
()-
1
);
}
}
String
step1
=
StrUtil
.
replace
(
dataJsonStr
,
"\\"
,
"\\\\"
);
boolean
kingbase
=
"kingbase"
.
equalsIgnoreCase
(
MapUtil
.
getStr
(
dbInfoMap
,
"db_dialect"
,
"mysql"
));
String
step2
=
StrUtil
.
replace
(
step1
,
"'"
,
"\\'"
);
// Kingbase(PostgreSQL) 字符串内的单引号转义为 '';MySQL 使用 \'
return
Tuple6
.
of
(
table
,
type
,
pkColumns
.
toString
(),
pkValues
.
toString
().
replace
(
"'"
,
""
),
step2
,
ts
);
String
escaped
=
kingbase
?
StrUtil
.
replace
(
dataJsonStr
,
"'"
,
"''"
)
:
StrUtil
.
replace
(
StrUtil
.
replace
(
dataJsonStr
,
"\\"
,
"\\\\"
),
"'"
,
"\\'"
);
return
Tuple6
.
of
(
table
,
type
,
pkColumns
.
toString
(),
pkValues
.
toString
().
replace
(
"'"
,
""
),
escaped
,
ts
);
}
}
private
static
String
tranferInsertSql
(
String
table
,
JSONObject
dataObj
,
JSONObject
mysqlType
)
{
private
static
String
tranferInsertSql
(
String
table
,
JSONObject
dataObj
,
JSONObject
mysqlType
)
{
...
@@ -131,6 +141,34 @@ public class MysqlDataTransferFunction extends ProcessFunction<JSONObject, Tuple
...
@@ -131,6 +141,34 @@ public class MysqlDataTransferFunction extends ProcessFunction<JSONObject, Tuple
return
sb
.
toString
();
return
sb
.
toString
();
}
}
/**
* Kingbase(PostgreSQL 系)写法:INSERT ... ON CONFLICT (pk) DO UPDATE
* 使用双引号包裹标识符(这里为兼容直接去掉特殊符号,采用无引号写法),单引号内转义用 '' 而非 \'
*/
private
static
String
transferUpsertSqlKingbase
(
String
table
,
JSONObject
dataObj
,
JSONObject
mysqlType
,
Set
<
String
>
pkNameSet
)
{
Set
<
String
>
columnSet
=
mysqlType
.
keySet
();
StringBuilder
colBuilder
=
new
StringBuilder
();
List
<
String
>
valueList
=
new
ArrayList
<>();
List
<
String
>
updateList
=
new
ArrayList
<>();
for
(
String
s
:
columnSet
)
{
colBuilder
.
append
(
s
).
append
(
","
);
String
val
=
getValueStringKingbase
(
dataObj
,
s
,
mysqlType
.
getString
(
s
));
valueList
.
add
(
val
);
updateList
.
add
(
s
+
" = EXCLUDED."
+
s
);
}
colBuilder
.
setLength
(
colBuilder
.
length
()
-
1
);
String
valueString
=
String
.
join
(
","
,
valueList
);
String
updateString
=
String
.
join
(
","
,
updateList
);
List
<
String
>
pkList
=
new
ArrayList
<>(
pkNameSet
);
String
pkString
=
String
.
join
(
","
,
pkList
);
// 若没有主键,退化为普通 INSERT(重复会报错,由下游表约束保证)
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) ON CONFLICT (%s) DO UPDATE SET %s"
,
table
,
colBuilder
.
toString
(),
valueString
,
pkString
,
updateString
);
}
private
static
String
transferDeleteSql
(
String
table
,
JSONObject
dataObj
,
JSONObject
mysqlType
,
Set
<
String
>
pkNameSet
)
{
private
static
String
transferDeleteSql
(
String
table
,
JSONObject
dataObj
,
JSONObject
mysqlType
,
Set
<
String
>
pkNameSet
)
{
StringBuilder
whereClauseBuilder
=
new
StringBuilder
();
StringBuilder
whereClauseBuilder
=
new
StringBuilder
();
for
(
String
pk
:
pkNameSet
)
{
for
(
String
pk
:
pkNameSet
)
{
...
@@ -169,4 +207,27 @@ public class MysqlDataTransferFunction extends ProcessFunction<JSONObject, Tuple
...
@@ -169,4 +207,27 @@ public class MysqlDataTransferFunction extends ProcessFunction<JSONObject, Tuple
return
dataObj
.
getString
(
columnKey
);
return
dataObj
.
getString
(
columnKey
);
}
}
/**
* Kingbase(PostgreSQL 系) 取值:字符串加单引号,内部单引号转义为 ''(而非 \')
*/
private
static
String
getValueStringKingbase
(
JSONObject
dataObj
,
String
columnKey
,
String
mysqlType
)
{
if
(
null
==
dataObj
.
get
(
columnKey
))
{
return
"null"
;
}
String
upperCase
=
mysqlType
.
toUpperCase
();
if
(
STR_SQL_TYPE
.
containsKey
(
upperCase
))
{
return
"'"
+
StrUtil
.
replace
(
dataObj
.
getString
(
columnKey
),
"'"
,
"''"
)
+
"'"
;
}
if
(
"DATE"
.
equals
(
upperCase
)
||
"DATETIME"
.
equals
(
upperCase
))
{
Date
d
=
dataObj
.
getDate
(
columnKey
);
if
(
d
==
null
)
{
return
""
;
}
LocalDateTime
dateTime
=
LocalDateTime
.
ofInstant
(
d
.
toInstant
(),
ZoneId
.
systemDefault
());
String
date
=
"DATETIME"
.
equals
(
upperCase
)
?
DATETIME_FORMAT
.
format
(
dateTime
)
:
DATE_FORMAT
.
format
(
dateTime
);
return
String
.
format
(
"'%s'"
,
date
);
}
return
dataObj
.
getString
(
columnKey
);
}
}
}
src/main/java/com/dsk/flink/dsc/common/sink/
Mysql
DataSlideSink.java
→
src/main/java/com/dsk/flink/dsc/common/sink/
Db
DataSlideSink.java
View file @
418dd22d
...
@@ -3,6 +3,7 @@ package com.dsk.flink.dsc.common.sink;
...
@@ -3,6 +3,7 @@ package com.dsk.flink.dsc.common.sink;
import
cn.hutool.core.lang.Snowflake
;
import
cn.hutool.core.lang.Snowflake
;
import
cn.hutool.core.util.IdUtil
;
import
cn.hutool.core.util.IdUtil
;
import
cn.hutool.core.util.RandomUtil
;
import
cn.hutool.core.util.RandomUtil
;
import
cn.hutool.core.util.StrUtil
;
import
cn.hutool.db.DbUtil
;
import
cn.hutool.db.DbUtil
;
import
cn.hutool.db.sql.SqlExecutor
;
import
cn.hutool.db.sql.SqlExecutor
;
import
com.alibaba.druid.pool.DruidDataSource
;
import
com.alibaba.druid.pool.DruidDataSource
;
...
@@ -22,17 +23,17 @@ import java.util.concurrent.LinkedBlockingDeque;
...
@@ -22,17 +23,17 @@ 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
;
public
class
Mysql
DataSlideSink
extends
RichSinkFunction
<
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>>
{
public
class
Db
DataSlideSink
extends
RichSinkFunction
<
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>>
{
static
Logger
logger
=
LoggerFactory
.
getLogger
(
Mysql
DataSlideSink
.
class
);
static
Logger
logger
=
LoggerFactory
.
getLogger
(
Db
DataSlideSink
.
class
);
EnvProperties
envProps
;
EnvProperties
envProps
;
private
transient
ExecutorService
executorService
;
private
transient
ExecutorService
executorService
;
private
transient
DruidDataSource
dataSource
;
private
transient
DruidDataSource
dataSource
;
private
static
final
int
MAX_RETRIES
=
3
;
// 最大重试次数
private
static
final
int
MAX_RETRIES
=
3
;
// 最大重试次数
private
static
final
int
RETRY_DELAY_MS
=
100
;
// 重试间隔时间
private
static
final
int
RETRY_DELAY_MS
=
100
;
// 重试间隔时间
private
static
final
String
SQL
=
"INSERT INTO dsc_cdc_log (
`table`
,op_type,pk_columns,pk_values,data_json,cdc_ts) values (?,?,?,?,?,?)"
;
private
static
final
String
SQL
=
"INSERT INTO dsc_cdc_log (
table
,op_type,pk_columns,pk_values,data_json,cdc_ts) values (?,?,?,?,?,?)"
;
public
Mysql
DataSlideSink
(
EnvProperties
envProps
)
{
public
Db
DataSlideSink
(
EnvProperties
envProps
)
{
this
.
envProps
=
envProps
;
this
.
envProps
=
envProps
;
}
}
...
@@ -43,7 +44,13 @@ public class MysqlDataSlideSink extends RichSinkFunction<Tuple6<String,String,St
...
@@ -43,7 +44,13 @@ public class MysqlDataSlideSink extends RichSinkFunction<Tuple6<String,String,St
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
();
dataSource
.
setDriverClassName
(
"com.mysql.cj.jdbc.Driver"
);
// 驱动类名优先取配置 db_driver,未配置则按方言推断,仍无则默认 MySQL
String
driverClass
=
envProps
.
getDb_driver
();
if
(
StrUtil
.
isBlank
(
driverClass
))
{
driverClass
=
"kingbase"
.
equalsIgnoreCase
(
envProps
.
getDb_dialect
())
?
"com.kingbase8.Driver"
:
"com.mysql.cj.jdbc.Driver"
;
}
dataSource
.
setDriverClassName
(
driverClass
);
dataSource
.
setUsername
(
envProps
.
getDb_username
());
dataSource
.
setUsername
(
envProps
.
getDb_username
());
dataSource
.
setPassword
(
envProps
.
getDb_password
());
dataSource
.
setPassword
(
envProps
.
getDb_password
());
dataSource
.
setUrl
(
configTidbUrl
);
dataSource
.
setUrl
(
configTidbUrl
);
...
@@ -100,7 +107,7 @@ public class MysqlDataSlideSink extends RichSinkFunction<Tuple6<String,String,St
...
@@ -100,7 +107,7 @@ public class MysqlDataSlideSink extends RichSinkFunction<Tuple6<String,String,St
private
void
writeErrLogDb
(
SqlErrorLog
errorLog
)
{
private
void
writeErrLogDb
(
SqlErrorLog
errorLog
)
{
Snowflake
snowflake
=
IdUtil
.
getSnowflake
(
RandomUtil
.
randomInt
(
31
),
RandomUtil
.
randomInt
(
31
));
Snowflake
snowflake
=
IdUtil
.
getSnowflake
(
RandomUtil
.
randomInt
(
31
),
RandomUtil
.
randomInt
(
31
));
String
sql
=
"insert dsc_err_log (id,error_time, error_sql, error_msg) values (?, ?, ?, ?)"
;
String
sql
=
"insert
into
dsc_err_log (id,error_time, error_sql, error_msg) values (?, ?, ?, ?)"
;
Connection
conn
=
null
;
Connection
conn
=
null
;
PreparedStatement
pt
=
null
;
PreparedStatement
pt
=
null
;
try
{
try
{
...
...
src/main/java/com/dsk/flink/dsc/common/sink/
Mysql
DataTransferSink.java
→
src/main/java/com/dsk/flink/dsc/common/sink/
Db
DataTransferSink.java
View file @
418dd22d
...
@@ -3,6 +3,7 @@ package com.dsk.flink.dsc.common.sink;
...
@@ -3,6 +3,7 @@ package com.dsk.flink.dsc.common.sink;
import
cn.hutool.core.lang.Snowflake
;
import
cn.hutool.core.lang.Snowflake
;
import
cn.hutool.core.util.IdUtil
;
import
cn.hutool.core.util.IdUtil
;
import
cn.hutool.core.util.RandomUtil
;
import
cn.hutool.core.util.RandomUtil
;
import
cn.hutool.core.util.StrUtil
;
import
cn.hutool.db.DbUtil
;
import
cn.hutool.db.DbUtil
;
import
cn.hutool.db.sql.SqlExecutor
;
import
cn.hutool.db.sql.SqlExecutor
;
import
com.alibaba.druid.pool.DruidDataSource
;
import
com.alibaba.druid.pool.DruidDataSource
;
...
@@ -21,16 +22,16 @@ import java.util.concurrent.LinkedBlockingDeque;
...
@@ -21,16 +22,16 @@ 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
;
public
class
Mysql
DataTransferSink
extends
RichSinkFunction
<
String
>
{
public
class
Db
DataTransferSink
extends
RichSinkFunction
<
String
>
{
static
Logger
logger
=
LoggerFactory
.
getLogger
(
Mysql
DataTransferSink
.
class
);
static
Logger
logger
=
LoggerFactory
.
getLogger
(
Db
DataTransferSink
.
class
);
EnvProperties
envProps
;
EnvProperties
envProps
;
private
transient
ExecutorService
executorService
;
private
transient
ExecutorService
executorService
;
private
transient
DruidDataSource
dataSource
;
private
transient
DruidDataSource
dataSource
;
private
static
final
int
MAX_RETRIES
=
3
;
// 最大重试次数
private
static
final
int
MAX_RETRIES
=
3
;
// 最大重试次数
private
static
final
int
RETRY_DELAY_MS
=
100
;
// 重试间隔时间
private
static
final
int
RETRY_DELAY_MS
=
100
;
// 重试间隔时间
public
Mysql
DataTransferSink
(
EnvProperties
envProps
)
{
public
Db
DataTransferSink
(
EnvProperties
envProps
)
{
this
.
envProps
=
envProps
;
this
.
envProps
=
envProps
;
}
}
...
@@ -41,7 +42,13 @@ public class MysqlDataTransferSink extends RichSinkFunction<String> {
...
@@ -41,7 +42,13 @@ public class MysqlDataTransferSink 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
();
dataSource
.
setDriverClassName
(
"com.mysql.cj.jdbc.Driver"
);
// 驱动类名优先取配置 db_driver,未配置则按方言推断,仍无则默认 MySQL
String
driverClass
=
envProps
.
getDb_driver
();
if
(
StrUtil
.
isBlank
(
driverClass
))
{
driverClass
=
"kingbase"
.
equalsIgnoreCase
(
envProps
.
getDb_dialect
())
?
"com.kingbase8.Driver"
:
"com.mysql.cj.jdbc.Driver"
;
}
dataSource
.
setDriverClassName
(
driverClass
);
dataSource
.
setUsername
(
envProps
.
getDb_username
());
dataSource
.
setUsername
(
envProps
.
getDb_username
());
dataSource
.
setPassword
(
envProps
.
getDb_password
());
dataSource
.
setPassword
(
envProps
.
getDb_password
());
dataSource
.
setUrl
(
configTidbUrl
);
dataSource
.
setUrl
(
configTidbUrl
);
...
@@ -117,7 +124,7 @@ public class MysqlDataTransferSink extends RichSinkFunction<String> {
...
@@ -117,7 +124,7 @@ public class MysqlDataTransferSink extends RichSinkFunction<String> {
private
void
writeErrLogDb
(
SqlErrorLog
errorLog
)
{
private
void
writeErrLogDb
(
SqlErrorLog
errorLog
)
{
Snowflake
snowflake
=
IdUtil
.
getSnowflake
(
RandomUtil
.
randomInt
(
31
),
RandomUtil
.
randomInt
(
31
));
Snowflake
snowflake
=
IdUtil
.
getSnowflake
(
RandomUtil
.
randomInt
(
31
),
RandomUtil
.
randomInt
(
31
));
String
sql
=
"insert dsc_err_log (id,error_time, error_sql, error_msg) values (?, ?, ?, ?)"
;
String
sql
=
"insert
into
dsc_err_log (id,error_time, error_sql, error_msg) values (?, ?, ?, ?)"
;
Connection
conn
=
null
;
Connection
conn
=
null
;
PreparedStatement
pt
=
null
;
PreparedStatement
pt
=
null
;
try
{
try
{
...
...
src/main/java/com/dsk/flink/dsc/common/sink/
Mysql
DataTransferSinkBatch.java
→
src/main/java/com/dsk/flink/dsc/common/sink/
Db
DataTransferSinkBatch.java
View file @
418dd22d
...
@@ -2,6 +2,7 @@ package com.dsk.flink.dsc.common.sink;
...
@@ -2,6 +2,7 @@ package com.dsk.flink.dsc.common.sink;
import
cn.hutool.core.collection.CollUtil
;
import
cn.hutool.core.collection.CollUtil
;
import
cn.hutool.core.lang.Snowflake
;
import
cn.hutool.core.lang.Snowflake
;
import
cn.hutool.core.util.StrUtil
;
import
cn.hutool.core.util.IdUtil
;
import
cn.hutool.core.util.IdUtil
;
import
cn.hutool.core.util.RandomUtil
;
import
cn.hutool.core.util.RandomUtil
;
import
com.alibaba.druid.pool.DruidDataSource
;
import
com.alibaba.druid.pool.DruidDataSource
;
...
@@ -25,9 +26,9 @@ import java.util.concurrent.atomic.AtomicBoolean;
...
@@ -25,9 +26,9 @@ import java.util.concurrent.atomic.AtomicBoolean;
* @author lww
* @author lww
* @date 2025-01-14
* @date 2025-01-14
*/
*/
public
class
Mysql
DataTransferSinkBatch
extends
RichSinkFunction
<
String
>
{
public
class
Db
DataTransferSinkBatch
extends
RichSinkFunction
<
String
>
{
static
Logger
logger
=
LoggerFactory
.
getLogger
(
MysqlDataTransferSink
.
class
);
static
Logger
logger
=
LoggerFactory
.
getLogger
(
DbDataTransferSinkBatch
.
class
);
EnvProperties
envProps
;
EnvProperties
envProps
;
private
transient
ExecutorService
executorService
;
private
transient
ExecutorService
executorService
;
private
transient
DruidDataSource
dataSource
;
private
transient
DruidDataSource
dataSource
;
...
@@ -40,7 +41,7 @@ public class MysqlDataTransferSinkBatch extends RichSinkFunction<String> {
...
@@ -40,7 +41,7 @@ public class MysqlDataTransferSinkBatch extends RichSinkFunction<String> {
private
AtomicBoolean
flushing
=
new
AtomicBoolean
(
false
);
private
AtomicBoolean
flushing
=
new
AtomicBoolean
(
false
);
private
int
subtaskIndex
;
private
int
subtaskIndex
;
public
Mysql
DataTransferSinkBatch
(
EnvProperties
envProps
)
{
public
Db
DataTransferSinkBatch
(
EnvProperties
envProps
)
{
this
.
envProps
=
envProps
;
this
.
envProps
=
envProps
;
}
}
...
@@ -51,7 +52,13 @@ public class MysqlDataTransferSinkBatch extends RichSinkFunction<String> {
...
@@ -51,7 +52,13 @@ public class MysqlDataTransferSinkBatch 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
();
dataSource
.
setDriverClassName
(
"com.mysql.cj.jdbc.Driver"
);
// 驱动类名优先取配置 db_driver,未配置则按方言推断,仍无则默认 MySQL
String
driverClass
=
envProps
.
getDb_driver
();
if
(
StrUtil
.
isBlank
(
driverClass
))
{
driverClass
=
"kingbase"
.
equalsIgnoreCase
(
envProps
.
getDb_dialect
())
?
"com.kingbase8.Driver"
:
"com.mysql.cj.jdbc.Driver"
;
}
dataSource
.
setDriverClassName
(
driverClass
);
dataSource
.
setUsername
(
envProps
.
getDb_username
());
dataSource
.
setUsername
(
envProps
.
getDb_username
());
dataSource
.
setPassword
(
envProps
.
getDb_password
());
dataSource
.
setPassword
(
envProps
.
getDb_password
());
dataSource
.
setUrl
(
configTidbUrl
);
dataSource
.
setUrl
(
configTidbUrl
);
...
@@ -60,7 +67,10 @@ public class MysqlDataTransferSinkBatch extends RichSinkFunction<String> {
...
@@ -60,7 +67,10 @@ public class MysqlDataTransferSinkBatch extends RichSinkFunction<String> {
dataSource
.
setTestWhileIdle
(
true
);
dataSource
.
setTestWhileIdle
(
true
);
dataSource
.
setMaxWait
(
20000
);
dataSource
.
setMaxWait
(
20000
);
dataSource
.
setValidationQuery
(
"select 1"
);
dataSource
.
setValidationQuery
(
"select 1"
);
dataSource
.
setConnectionInitSqls
(
CollUtil
.
newArrayList
(
"SET SESSION TRANSACTION ISOLATION LEVEL READ COMMITTED"
));
// MySQL 专属的连接初始化 SQL,Kingbase 不兼容,按方言跳过
if
(!
"kingbase"
.
equalsIgnoreCase
(
envProps
.
getDb_dialect
()))
{
dataSource
.
setConnectionInitSqls
(
CollUtil
.
newArrayList
(
"SET SESSION TRANSACTION ISOLATION LEVEL READ COMMITTED"
));
}
scheduledExecutorService
=
new
ScheduledThreadPoolExecutor
(
1
);
scheduledExecutorService
=
new
ScheduledThreadPoolExecutor
(
1
);
scheduledExecutorService
.
scheduleAtFixedRate
(
this
::
flush
,
FLUSH_INTERVAL
,
FLUSH_INTERVAL
,
TimeUnit
.
MILLISECONDS
);
scheduledExecutorService
.
scheduleAtFixedRate
(
this
::
flush
,
FLUSH_INTERVAL
,
FLUSH_INTERVAL
,
TimeUnit
.
MILLISECONDS
);
logger
.
info
(
"Subtask {} initialized with batch size {}"
,
subtaskIndex
,
BATCH_SIZE
);
logger
.
info
(
"Subtask {} initialized with batch size {}"
,
subtaskIndex
,
BATCH_SIZE
);
...
...
src/main/java/com/dsk/flink/dsc/sync/SyncCustomerDataSource.java
View file @
418dd22d
...
@@ -3,9 +3,9 @@ package com.dsk.flink.dsc.sync;
...
@@ -3,9 +3,9 @@ package com.dsk.flink.dsc.sync;
import
cn.hutool.core.map.MapUtil
;
import
cn.hutool.core.map.MapUtil
;
import
cn.hutool.core.util.StrUtil
;
import
cn.hutool.core.util.StrUtil
;
import
com.alibaba.fastjson.JSONObject
;
import
com.alibaba.fastjson.JSONObject
;
import
com.dsk.flink.dsc.common.function.
Mysql
DataTransferFunction
;
import
com.dsk.flink.dsc.common.function.
Db
DataTransferFunction
;
import
com.dsk.flink.dsc.common.sink.
Mysql
DataSlideSink
;
import
com.dsk.flink.dsc.common.sink.
Db
DataSlideSink
;
import
com.dsk.flink.dsc.common.sink.
Mysql
DataTransferSink
;
import
com.dsk.flink.dsc.common.sink.
Db
DataTransferSink
;
import
com.dsk.flink.dsc.utils.EnvProperties
;
import
com.dsk.flink.dsc.utils.EnvProperties
;
import
com.dsk.flink.dsc.utils.EnvPropertiesUtil
;
import
com.dsk.flink.dsc.utils.EnvPropertiesUtil
;
import
com.dsk.flink.dsc.utils.EtlUtils
;
import
com.dsk.flink.dsc.utils.EtlUtils
;
...
@@ -108,7 +108,7 @@ public class SyncCustomerDataSource {
...
@@ -108,7 +108,7 @@ public class SyncCustomerDataSource {
OutputTag
<
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>>
logSlideTag
=
new
OutputTag
<
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>>(
"log_slide"
)
{};
OutputTag
<
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>>
logSlideTag
=
new
OutputTag
<
Tuple6
<
String
,
String
,
String
,
String
,
String
,
Long
>>(
"log_slide"
)
{};
SingleOutputStreamOperator
<
Tuple3
<
String
,
String
,
Long
>>
slide
=
tsGroupStream
SingleOutputStreamOperator
<
Tuple3
<
String
,
String
,
Long
>>
slide
=
tsGroupStream
.
process
(
new
Mysql
DataTransferFunction
(
envProps
,
logSlideTag
))
.
process
(
new
Db
DataTransferFunction
(
envProps
,
logSlideTag
))
.
name
(
"dsc-sql"
)
.
name
(
"dsc-sql"
)
.
uid
(
"dsc-sql"
);
.
uid
(
"dsc-sql"
);
...
@@ -134,7 +134,7 @@ public class SyncCustomerDataSource {
...
@@ -134,7 +134,7 @@ public class SyncCustomerDataSource {
.
name
(
"dsc-max"
)
.
name
(
"dsc-max"
)
.
uid
(
"dsc-max"
);
.
uid
(
"dsc-max"
);
groupWindowSqlResultStream
.
addSink
(
new
Mysql
DataTransferSink
(
envProps
))
groupWindowSqlResultStream
.
addSink
(
new
Db
DataTransferSink
(
envProps
))
.
name
(
"dsc-sink"
)
.
name
(
"dsc-sink"
)
.
uid
(
"dsc-sink"
);
.
uid
(
"dsc-sink"
);
...
@@ -164,7 +164,7 @@ public class SyncCustomerDataSource {
...
@@ -164,7 +164,7 @@ public class SyncCustomerDataSource {
)).uid("dsc-log")
)).uid("dsc-log")
.name("dsc-log");*/
.name("dsc-log");*/
sideOutput
.
addSink
(
new
Mysql
DataSlideSink
(
envProps
)).
uid
(
"dsc-log"
)
sideOutput
.
addSink
(
new
Db
DataSlideSink
(
envProps
)).
uid
(
"dsc-log"
)
.
name
(
"dsc-log"
);
.
name
(
"dsc-log"
);
env
.
execute
(
"dsc-client"
);
env
.
execute
(
"dsc-client"
);
...
...
src/main/java/com/dsk/flink/dsc/utils/EnvProperties.java
View file @
418dd22d
...
@@ -13,6 +13,12 @@ public class EnvProperties extends Properties {
...
@@ -13,6 +13,12 @@ public class EnvProperties extends Properties {
// #连接
// #连接
String
db_url
;
String
db_url
;
// #数据库驱动类全名(如 com.mysql.cj.jdbc.Driver / com.kingbase8.Driver)
String
db_driver
;
// #数据库方言:mysql / kingbase,决定 SQL 生成与转义规则
String
db_dialect
;
// #建设库TIDB库
// #建设库TIDB库
String
db_host
;
String
db_host
;
String
db_port
;
String
db_port
;
...
@@ -555,4 +561,20 @@ public class EnvProperties extends Properties {
...
@@ -555,4 +561,20 @@ public class EnvProperties extends Properties {
public
void
setKafka_topic
(
String
kafka_topic
)
{
public
void
setKafka_topic
(
String
kafka_topic
)
{
this
.
kafka_topic
=
kafka_topic
;
this
.
kafka_topic
=
kafka_topic
;
}
}
public
String
getDb_driver
()
{
return
db_driver
==
null
?
this
.
getProperty
(
"db_driver"
)
:
db_driver
;
}
public
void
setDb_driver
(
String
db_driver
)
{
this
.
db_driver
=
db_driver
;
}
public
String
getDb_dialect
()
{
return
db_dialect
==
null
?
this
.
getProperty
(
"db_dialect"
)
:
db_dialect
;
}
public
void
setDb_dialect
(
String
db_dialect
)
{
this
.
db_dialect
=
db_dialect
;
}
}
}
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