Skip to content
Projects
Groups
Snippets
Help
Loading...
Sign in
Toggle navigation
F
flink-monitor
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
赵建伟
flink-monitor
Commits
27aa4705
You need to sign in or sign up before continuing.
Commit
27aa4705
authored
Mar 25, 2020
by
赵建伟
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
update codes
parent
c47bbeaf
Show whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
14 additions
and
30 deletions
+14
-30
PortraitMonitorMain.java
src/main/java/com/gmei/data/monitor/PortraitMonitorMain.java
+0
-17
PortraitMonitorMainAll.java
...in/java/com/gmei/data/monitor/PortraitMonitorMainAll.java
+14
-13
No files found.
src/main/java/com/gmei/data/monitor/PortraitMonitorMain.java
View file @
27aa4705
package
com
.
gmei
.
data
.
monitor
;
package
com
.
gmei
.
data
.
monitor
;
import
com.gmei.data.monitor.operator.PortraitMonitorErrOperator
;
import
com.gmei.data.monitor.operator.PortraitMonitorShdOperator
;
import
com.gmei.data.monitor.operator.PortraitMonitorShdOperator
;
import
com.gmei.data.monitor.operator.PortraitMonitorSucOperator
;
import
com.gmei.data.monitor.operator.PortraitMonitorSucOperator
;
import
com.gmei.data.monitor.source.PortraitKafkaSource
;
import
com.gmei.data.monitor.source.PortraitKafkaSource
;
...
@@ -13,8 +12,6 @@ import org.apache.flink.streaming.api.datastream.DataStream;
...
@@ -13,8 +12,6 @@ import org.apache.flink.streaming.api.datastream.DataStream;
import
org.apache.flink.streaming.api.environment.CheckpointConfig
;
import
org.apache.flink.streaming.api.environment.CheckpointConfig
;
import
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
;
import
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
;
import
java.text.SimpleDateFormat
;
/**
/**
* @ClassName PortraitMonitorMain
* @ClassName PortraitMonitorMain
* @Description: 画像打点实时监控主入口
* @Description: 画像打点实时监控主入口
...
@@ -33,7 +30,6 @@ public class PortraitMonitorMain {
...
@@ -33,7 +30,6 @@ public class PortraitMonitorMain {
String
maidianInTopic
=
parameterTool
.
get
(
"maidianInTopic"
,
"test11"
);
String
maidianInTopic
=
parameterTool
.
get
(
"maidianInTopic"
,
"test11"
);
String
backendInTopic
=
parameterTool
.
get
(
"backendInTopic"
,
"test12"
);
String
backendInTopic
=
parameterTool
.
get
(
"backendInTopic"
,
"test12"
);
String
portraitSucInTopic
=
parameterTool
.
get
(
"portraitSucInTopic"
,
"test13"
);
String
portraitSucInTopic
=
parameterTool
.
get
(
"portraitSucInTopic"
,
"test13"
);
String
portraitErrGroupId
=
parameterTool
.
get
(
"portraitErrGroupId"
,
"flink_monitor_err"
);
String
portraitShdGroupId
=
parameterTool
.
get
(
"portraitShdGroupId"
,
"flink_monitor_shd"
);
String
portraitShdGroupId
=
parameterTool
.
get
(
"portraitShdGroupId"
,
"flink_monitor_shd"
);
String
portraitSucGroupId
=
parameterTool
.
get
(
"portraitSucGroupId"
,
"flink_monitor_suc"
);
String
portraitSucGroupId
=
parameterTool
.
get
(
"portraitSucGroupId"
,
"flink_monitor_suc"
);
Integer
windowSize
=
parameterTool
.
getInt
(
"windowSize"
,
60
);
Integer
windowSize
=
parameterTool
.
getInt
(
"windowSize"
,
60
);
...
@@ -57,18 +53,6 @@ public class PortraitMonitorMain {
...
@@ -57,18 +53,6 @@ public class PortraitMonitorMain {
CheckpointConfig
config
=
env
.
getCheckpointConfig
();
CheckpointConfig
config
=
env
.
getCheckpointConfig
();
config
.
enableExternalizedCheckpoints
(
CheckpointConfig
.
ExternalizedCheckpointCleanup
.
RETAIN_ON_CANCELLATION
);
config
.
enableExternalizedCheckpoints
(
CheckpointConfig
.
ExternalizedCheckpointCleanup
.
RETAIN_ON_CANCELLATION
);
// 获取数据源
DataStream
portraitErrDataStream
=
new
PortraitKafkaSource
(
env
,
inBrokers
,
maidianInTopic
,
backendInTopic
,
portraitErrGroupId
,
batchSize
,
isStartFromEarliest
,
isStartFromLatest
,
startTime
).
getInstance
();
DataStream
portraitShdDataStream
=
new
PortraitKafkaSource
(
DataStream
portraitShdDataStream
=
new
PortraitKafkaSource
(
env
,
env
,
inBrokers
,
inBrokers
,
...
@@ -92,7 +76,6 @@ public class PortraitMonitorMain {
...
@@ -92,7 +76,6 @@ public class PortraitMonitorMain {
).
getInstance
();
).
getInstance
();
// 执行处理核心逻辑
// 执行处理核心逻辑
new
PortraitMonitorErrOperator
(
portraitErrDataStream
,
outJdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
).
run
();
new
PortraitMonitorShdOperator
(
portraitShdDataStream
,
windowSize
,
slideSize
,
outJdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
).
run
();
new
PortraitMonitorShdOperator
(
portraitShdDataStream
,
windowSize
,
slideSize
,
outJdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
).
run
();
new
PortraitMonitorSucOperator
(
portraitSucDataStream
,
windowSize
,
slideSize
,
outJdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
).
run
();
new
PortraitMonitorSucOperator
(
portraitSucDataStream
,
windowSize
,
slideSize
,
outJdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
).
run
();
...
...
src/main/java/com/gmei/data/monitor/PortraitMonitorMain
ShdSuc
.java
→
src/main/java/com/gmei/data/monitor/PortraitMonitorMain
All
.java
View file @
27aa4705
package
com
.
gmei
.
data
.
monitor
;
package
com
.
gmei
.
data
.
monitor
;
import
com.gmei.data.monitor.operator.PortraitMonitorErrOperator
;
import
com.gmei.data.monitor.operator.PortraitMonitorShdOperator
;
import
com.gmei.data.monitor.operator.PortraitMonitorShdOperator
;
import
com.gmei.data.monitor.operator.PortraitMonitorSucOperator
;
import
com.gmei.data.monitor.operator.PortraitMonitorSucOperator
;
import
com.gmei.data.monitor.source.PortraitKafkaSource
;
import
com.gmei.data.monitor.source.PortraitKafkaSource
;
...
@@ -20,7 +21,7 @@ import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
...
@@ -20,7 +21,7 @@ import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
* @Date 2020/3/18
* @Date 2020/3/18
* @Version V1.0
* @Version V1.0
**/
**/
public
class
PortraitMonitorMain
ShdSuc
{
public
class
PortraitMonitorMain
All
{
public
static
void
main
(
String
[]
args
)
throws
Exception
{
public
static
void
main
(
String
[]
args
)
throws
Exception
{
// 获取运行参数
// 获取运行参数
...
@@ -55,17 +56,17 @@ public class PortraitMonitorMainShdSuc {
...
@@ -55,17 +56,17 @@ public class PortraitMonitorMainShdSuc {
config
.
enableExternalizedCheckpoints
(
CheckpointConfig
.
ExternalizedCheckpointCleanup
.
RETAIN_ON_CANCELLATION
);
config
.
enableExternalizedCheckpoints
(
CheckpointConfig
.
ExternalizedCheckpointCleanup
.
RETAIN_ON_CANCELLATION
);
// 获取数据源
// 获取数据源
//
DataStream portraitErrDataStream = new PortraitKafkaSource(
DataStream
portraitErrDataStream
=
new
PortraitKafkaSource
(
//
env,
env
,
//
inBrokers,
inBrokers
,
//
maidianInTopic,
maidianInTopic
,
//
backendInTopic,
backendInTopic
,
//
portraitErrGroupId,
portraitErrGroupId
,
//
batchSize,
batchSize
,
//
isStartFromEarliest,
isStartFromEarliest
,
//
isStartFromLatest,
isStartFromLatest
,
//
startTime
startTime
//
).getInstance();
).
getInstance
();
DataStream
portraitShdDataStream
=
new
PortraitKafkaSource
(
DataStream
portraitShdDataStream
=
new
PortraitKafkaSource
(
env
,
env
,
inBrokers
,
inBrokers
,
...
@@ -89,7 +90,7 @@ public class PortraitMonitorMainShdSuc {
...
@@ -89,7 +90,7 @@ public class PortraitMonitorMainShdSuc {
).
getInstance
();
).
getInstance
();
// 执行处理核心逻辑
// 执行处理核心逻辑
//
new PortraitMonitorErrOperator(portraitErrDataStream,outJdbcUrl,maxRetry,retryInteral,parallelism).run();
new
PortraitMonitorErrOperator
(
portraitErrDataStream
,
outJdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
).
run
();
new
PortraitMonitorShdOperator
(
portraitShdDataStream
,
windowSize
,
slideSize
,
outJdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
).
run
();
new
PortraitMonitorShdOperator
(
portraitShdDataStream
,
windowSize
,
slideSize
,
outJdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
).
run
();
new
PortraitMonitorSucOperator
(
portraitSucDataStream
,
windowSize
,
slideSize
,
outJdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
).
run
();
new
PortraitMonitorSucOperator
(
portraitSucDataStream
,
windowSize
,
slideSize
,
outJdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
).
run
();
...
...
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