Skip to content
Projects
Groups
Snippets
Help
Loading...
Sign in
Toggle navigation
C
ctr-estimate
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
赵建伟
ctr-estimate
Commits
3ba5e63c
Commit
3ba5e63c
authored
Apr 05, 2020
by
赵建伟
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
update codes
parent
6213c3f6
Hide whitespace changes
Inline
Side-by-side
Showing
10 changed files
with
245 additions
and
30 deletions
+245
-30
start.sh
bin/start.sh
+4
-3
start_clk.sh
bin/start_clk.sh
+4
-4
start_tag.sh
bin/start_tag.sh
+4
-4
DevCtrEstimateMainClk.java
src/main/java/com/gmei/data/ctr/DevCtrEstimateMainClk.java
+9
-11
DevCtrEstimateMainTag.java
src/main/java/com/gmei/data/ctr/DevCtrEstimateMainTag.java
+103
-0
ProdCtrEstimateMain.java
src/main/java/com/gmei/data/ctr/ProdCtrEstimateMain.java
+113
-0
ProdCtrEstimateMainClk.java
src/main/java/com/gmei/data/ctr/ProdCtrEstimateMainClk.java
+2
-2
ProdCtrEstimateMainTag.java
src/main/java/com/gmei/data/ctr/ProdCtrEstimateMainTag.java
+2
-2
TestCtrEstimateMainClk.java
src/main/java/com/gmei/data/ctr/TestCtrEstimateMainClk.java
+2
-2
TestCtrEstimateMainTag.java
src/main/java/com/gmei/data/ctr/TestCtrEstimateMainTag.java
+2
-2
No files found.
bin/start.sh
View file @
3ba5e63c
...
...
@@ -14,19 +14,20 @@ nohup $FLINK_HOME/bin/flink run \
-p
6
\
-yjm
1024
\
-ytm
2048
\
-c
com.gmei.data.ctr.ProdCtrEstimateMain
\
$JAR_DIR
/ctr-estimate-1.0-SNAPSHOT.jar
\
--inBrokers
'172.16.44.25:9092,172.16.44.31:9092,172.16.44.45:9092'
\
--batchSize
1000
\
--maidianInTopic
'gm-maidian-data'
\
--maidianInGroupId
'
ctr-estimate-flink
'
\
--maidianInGroupId
'
prod-ctr-estimate
'
\
--windowSize
5
\
--slideSize
5
\
--jdbcUrl
'jdbc:mysql://172.1
8.44.3:3306/jerry_test?user=root&password=5OqYM^zLwotJ3oSo
&autoReconnect=true&useSSL=false'
\
--jdbcUrl
'jdbc:mysql://172.1
6.40.170:4000/jerry_test?user=data_user&password=YPEzp78HQBuhByWPpefQu6X3D6hEPfD6
&autoReconnect=true&useSSL=false'
\
--maxRetry
3
\
--retryInteral
3000
\
--checkpointPath
'hdfs://bj-gmei-hdfs/user/data/flink/ctr-estimate/checkpoint'
\
--parallelism
6
\
--startTime
'2020-04-0
4 10:55
:00'
\
--startTime
'2020-04-0
5 00:00
:00'
\
>>
/data/log/ctr-estimate/ctr-estimate.out 2>&1 &
tail
-f
/data/log/ctr-estimate/ctr-estimate.out
...
...
bin/start_clk.sh
View file @
3ba5e63c
...
...
@@ -14,20 +14,20 @@ nohup $FLINK_HOME/bin/flink run \
-p
6
\
-yjm
1024
\
-ytm
2048
\
-c
com.gmei.data.ctr.CtrEstimateMainClk
\
-c
com.gmei.data.ctr.
Prod
CtrEstimateMainClk
\
$JAR_DIR
/ctr-estimate-1.0-SNAPSHOT.jar
\
--inBrokers
'172.16.44.25:9092,172.16.44.31:9092,172.16.44.45:9092'
\
--batchSize
1000
\
--maidianInTopic
'gm-maidian-data'
\
--maidianInGroupId
'
ctr-estimate-flink
-clk'
\
--maidianInGroupId
'
test-ctr-estimate
-clk'
\
--windowSize
5
\
--slideSize
5
\
--jdbcUrl
'jdbc:mysql://172.1
8.44.3:3306/jerry_test?user=root&password=5OqYM^zLwotJ3oSo
&autoReconnect=true&useSSL=false'
\
--jdbcUrl
'jdbc:mysql://172.1
6.40.170:4000/jerry_test?user=data_user&password=YPEzp78HQBuhByWPpefQu6X3D6hEPfD6
&autoReconnect=true&useSSL=false'
\
--maxRetry
3
\
--retryInteral
3000
\
--checkpointPath
'hdfs://bj-gmei-hdfs/user/data/flink/ctr-estimate-clk/checkpoint'
\
--parallelism
6
\
--startTime
'2020-04-0
4 15:23
:00'
\
--startTime
'2020-04-0
5 00:00
:00'
\
>>
/data/log/ctr-estimate/ctr-estimate-clk.out 2>&1 &
tail
-f
/data/log/ctr-estimate/ctr-estimate-clk.out
...
...
bin/start_tag.sh
View file @
3ba5e63c
...
...
@@ -14,20 +14,20 @@ nohup $FLINK_HOME/bin/flink run \
-p
6
\
-yjm
1024
\
-ytm
2048
\
-c
com.gmei.data.ctr.CtrEstimateMainTag
\
-c
com.gmei.data.ctr.
Prod
CtrEstimateMainTag
\
$JAR_DIR
/ctr-estimate-1.0-SNAPSHOT.jar
\
--inBrokers
'172.16.44.25:9092,172.16.44.31:9092,172.16.44.45:9092'
\
--batchSize
1000
\
--maidianInTopic
'gm-maidian-data'
\
--maidianInGroupId
'
ctr-estimate-flink
-tag'
\
--maidianInGroupId
'
test-ctr-estimate
-tag'
\
--windowSize
5
\
--slideSize
5
\
--jdbcUrl
'jdbc:mysql://172.1
8.44.3:3306/jerry_test?user=root&password=5OqYM^zLwotJ3oSo
&autoReconnect=true&useSSL=false'
\
--jdbcUrl
'jdbc:mysql://172.1
6.40.170:4000/jerry_test?user=data_user&password=YPEzp78HQBuhByWPpefQu6X3D6hEPfD6
&autoReconnect=true&useSSL=false'
\
--maxRetry
3
\
--retryInteral
3000
\
--checkpointPath
'hdfs://bj-gmei-hdfs/user/data/flink/ctr-estimate-tag/checkpoint'
\
--parallelism
6
\
--startTime
'2020-04-0
4
00:00:00'
\
--startTime
'2020-04-0
5
00:00:00'
\
>>
/data/log/ctr-estimate/ctr-estimate-tag.out 2>&1 &
tail
-f
/data/log/ctr-estimate/ctr-estimate-tag.out
...
...
src/main/java/com/gmei/data/ctr/
CtrEstimateMain
.java
→
src/main/java/com/gmei/data/ctr/
DevCtrEstimateMainClk
.java
View file @
3ba5e63c
package
com
.
gmei
.
data
.
ctr
;
import
com.gmei.data.ctr.operator.CtrEstimateClkOperator
;
import
com.gmei.data.ctr.operator.CtrEstimateTagOperator
;
import
com.gmei.data.ctr.source.MaidianKafkaSource
;
import
org.apache.flink.api.common.restartstrategy.RestartStrategies
;
import
org.apache.flink.api.java.utils.ParameterTool
;
...
...
@@ -11,13 +10,13 @@ import org.apache.flink.streaming.api.environment.CheckpointConfig;
import
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
;
/**
* @ClassName
CtrEstimateMain
* @ClassName
DevCtrEstimateMainClk
* @Description: CTR预估特征实时处理入口
* @Author apple
* @Date 2020/3/30
* @Version V1.0
**/
public
class
CtrEstimateMain
{
public
class
DevCtrEstimateMainClk
{
public
static
void
main
(
String
[]
args
)
throws
Exception
{
// 获取运行参数
...
...
@@ -25,7 +24,7 @@ public class CtrEstimateMain {
String
inBrokers
=
parameterTool
.
get
(
"inBrokers"
,
"test003:9092"
);
String
batchSize
=
parameterTool
.
get
(
"batchSize"
,
"1000"
);
String
maidianInTopic
=
parameterTool
.
get
(
"maidianInTopic"
,
"test11"
);
String
maidianInGroupId
=
parameterTool
.
get
(
"maidianInGroupId"
,
"ctr-estimate"
);
String
maidianInGroupId
=
parameterTool
.
get
(
"maidianInGroupId"
,
"ctr-estimate
-clk
"
);
Integer
windowSize
=
parameterTool
.
getInt
(
"windowSize"
,
60
);
Integer
slideSize
=
parameterTool
.
getInt
(
"slideSize"
,
60
);
String
jdbcUrl
=
parameterTool
.
get
(
"jdbcUrl"
,
...
...
@@ -52,11 +51,11 @@ public class CtrEstimateMain {
// 获得流处理环境对象
StreamExecutionEnvironment
env
=
StreamExecutionEnvironment
.
getExecutionEnvironment
();
//env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
env
.
enableCheckpointing
(
1000
);
env
.
setStateBackend
(
new
FsStateBackend
(
checkpointPath
));
env
.
setRestartStrategy
(
RestartStrategies
.
fixedDelayRestart
(
1
,
3000
));
CheckpointConfig
config
=
env
.
getCheckpointConfig
();
config
.
enableExternalizedCheckpoints
(
CheckpointConfig
.
ExternalizedCheckpointCleanup
.
RETAIN_ON_CANCELLATION
);
//
env.enableCheckpointing(1000);
//
env.setStateBackend(new FsStateBackend(checkpointPath));
//
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(1, 3000));
//
CheckpointConfig config = env.getCheckpointConfig();
//
config.enableExternalizedCheckpoints(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
DataStream
MaidianDataStream
=
new
MaidianKafkaSource
(
env
,
...
...
@@ -71,9 +70,8 @@ public class CtrEstimateMain {
// 执行处理核心逻辑
new
CtrEstimateClkOperator
(
MaidianDataStream
,
jdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
,
windowSize
,
slideSize
).
run
();
new
CtrEstimateTagOperator
(
MaidianDataStream
,
jdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
,
windowSize
,
slideSize
).
run
();
// 常驻执行
env
.
execute
(
"ctr-estimate"
);
env
.
execute
(
"ctr-estimate
-clk
"
);
}
}
src/main/java/com/gmei/data/ctr/DevCtrEstimateMainTag.java
0 → 100644
View file @
3ba5e63c
package
com
.
gmei
.
data
.
ctr
;
import
com.gmei.data.ctr.operator.CtrEstimateTagOperator
;
import
com.gmei.data.ctr.source.MaidianKafkaSource
;
import
org.apache.flink.api.common.restartstrategy.RestartStrategies
;
import
org.apache.flink.api.java.utils.ParameterTool
;
import
org.apache.flink.runtime.state.filesystem.FsStateBackend
;
import
org.apache.flink.streaming.api.datastream.DataStream
;
import
org.apache.flink.streaming.api.environment.CheckpointConfig
;
import
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
;
/**
* @ClassName DevCtrEstimateMainTag
* @Description: CTR预估特征实时处理入口
* @Author apple
* @Date 2020/3/30
* @Version V1.0
**/
public
class
DevCtrEstimateMainTag
{
public
static
void
main
(
String
[]
args
)
throws
Exception
{
// 获取运行参数
ParameterTool
parameterTool
=
ParameterTool
.
fromArgs
(
args
);
String
inBrokers
=
parameterTool
.
get
(
"inBrokers"
,
"test003:9092"
);
String
batchSize
=
parameterTool
.
get
(
"batchSize"
,
"1000"
);
String
maidianInTopic
=
parameterTool
.
get
(
"maidianInTopic"
,
"test11"
);
String
maidianInGroupId
=
parameterTool
.
get
(
"maidianInGroupId"
,
"ctr-estimate-tag"
);
Integer
windowSize
=
parameterTool
.
getInt
(
"windowSize"
,
5
);
Integer
slideSize
=
parameterTool
.
getInt
(
"slideSize"
,
5
);
String
jdbcUrl
=
parameterTool
.
get
(
"jdbcUrl"
,
"jdbc:mysql://172.18.44.3:3306/jerry_test?user=root&password=5OqYM^zLwotJ3oSo&autoReconnect=true&useSSL=false"
);
Integer
maxRetry
=
parameterTool
.
getInt
(
"maxRetry"
,
3
);
Long
retryInteral
=
parameterTool
.
getLong
(
"retryInteral"
,
3000
);
String
checkpointPath
=
parameterTool
.
get
(
"checkpointPath"
,
"hdfs://bj-gmei-hdfs/user/data/flink/ctr-estimate/checkpoint"
);
Boolean
isStartFromEarliest
=
parameterTool
.
getBoolean
(
"isStartFromEarliest"
,
true
);
Boolean
isStartFromLatest
=
parameterTool
.
getBoolean
(
"isStartFromLatest"
,
false
);
String
startTime
=
parameterTool
.
get
(
"startTime"
);
Integer
parallelism
=
parameterTool
.
getInt
(
"parallelism"
,
2
);
String
zxJdbcUrl
=
parameterTool
.
get
(
"zxJdbcUrl"
,
"jdbc:mysql://172.16.30.141:3306/zhengxing?characterEncoding=UTF-8&autoReconnect=true&useSSL=false"
);
String
zxUsername
=
parameterTool
.
get
(
"zxUsername"
,
"work"
);
String
zxPassword
=
parameterTool
.
get
(
"zxPassword"
,
"BJQaT9VzDcuPBqkd"
);
String
jerryJdbcUrl
=
parameterTool
.
get
(
"jerryJdbcUrl"
,
"jdbc:mysql://172.16.40.170:4000/jerry_test?characterEncoding=UTF-8&autoReconnect=true&useSSL=false"
);
String
jerryUsername
=
parameterTool
.
get
(
"jerryUsername"
,
"data_user"
);
String
jerryPassword
=
parameterTool
.
get
(
"jerryPassword"
,
"YPEzp78HQBuhByWPpefQu6X3D6hEPfD6"
);
System
.
out
.
println
(
"**********************************************************"
);
System
.
out
.
println
(
"*** inBrokers: "
+
inBrokers
);
System
.
out
.
println
(
"*** maidianInTopic: "
+
maidianInTopic
);
System
.
out
.
println
(
"*** maidianInGroupId: "
+
maidianInGroupId
);
System
.
out
.
println
(
"*** jdbcUrl: "
+
jdbcUrl
);
System
.
out
.
println
(
"*** checkpointPath: "
+
checkpointPath
);
System
.
out
.
println
(
"*** startTime: "
+
startTime
);
System
.
out
.
println
(
"*** windowSize: "
+
windowSize
);
System
.
out
.
println
(
"*** slideSize: "
+
slideSize
);
System
.
out
.
println
(
"*** zxJdbcUrl: "
+
zxJdbcUrl
);
System
.
out
.
println
(
"*** zxUsername: "
+
zxUsername
);
System
.
out
.
println
(
"*** zxPassword: "
+
zxPassword
);
System
.
out
.
println
(
"*** jerryJdbcUrl: "
+
jerryJdbcUrl
);
System
.
out
.
println
(
"*** jerryUsername: "
+
jerryUsername
);
System
.
out
.
println
(
"*** jerryPassword: "
+
jerryPassword
);
System
.
out
.
println
(
"**********************************************************"
);
// 获得流处理环境对象
StreamExecutionEnvironment
env
=
StreamExecutionEnvironment
.
getExecutionEnvironment
();
//env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
// env.enableCheckpointing(1000);
// env.setStateBackend(new FsStateBackend(checkpointPath));
// env.setRestartStrategy(RestartStrategies.fixedDelayRestart(1, 3000));
// CheckpointConfig config = env.getCheckpointConfig();
// config.enableExternalizedCheckpoints(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
DataStream
MaidianDataStream
=
new
MaidianKafkaSource
(
env
,
inBrokers
,
maidianInTopic
,
maidianInGroupId
,
batchSize
,
isStartFromEarliest
,
isStartFromLatest
,
startTime
).
getInstance
();
// 执行处理核心逻辑
new
CtrEstimateTagOperator
(
MaidianDataStream
,
jdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
,
windowSize
,
slideSize
,
zxJdbcUrl
,
zxUsername
,
zxPassword
,
jerryJdbcUrl
,
jerryUsername
,
jerryPassword
).
run
();
// 常驻执行
env
.
execute
(
"ctr-estimate-tag"
);
}
}
src/main/java/com/gmei/data/ctr/ProdCtrEstimateMain.java
0 → 100644
View file @
3ba5e63c
package
com
.
gmei
.
data
.
ctr
;
import
com.gmei.data.ctr.operator.CtrEstimateClkOperator
;
import
com.gmei.data.ctr.operator.CtrEstimateTagOperator
;
import
com.gmei.data.ctr.source.MaidianKafkaSource
;
import
org.apache.flink.api.common.restartstrategy.RestartStrategies
;
import
org.apache.flink.api.java.utils.ParameterTool
;
import
org.apache.flink.runtime.state.filesystem.FsStateBackend
;
import
org.apache.flink.streaming.api.datastream.DataStream
;
import
org.apache.flink.streaming.api.environment.CheckpointConfig
;
import
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
;
/**
* @ClassName ProdCtrEstimateMain
* @Description: CTR预估特征实时处理入口
* @Author apple
* @Date 2020/3/30
* @Version V1.0
**/
public
class
ProdCtrEstimateMain
{
public
static
void
main
(
String
[]
args
)
throws
Exception
{
// 获取运行参数
ParameterTool
parameterTool
=
ParameterTool
.
fromArgs
(
args
);
String
inBrokers
=
parameterTool
.
get
(
"inBrokers"
,
"test003:9092"
);
String
batchSize
=
parameterTool
.
get
(
"batchSize"
,
"1000"
);
String
maidianInTopic
=
parameterTool
.
get
(
"maidianInTopic"
,
"test11"
);
String
maidianInGroupId
=
parameterTool
.
get
(
"maidianInGroupId"
,
"ctr-estimate"
);
Integer
windowSize
=
parameterTool
.
getInt
(
"windowSize"
,
60
);
Integer
slideSize
=
parameterTool
.
getInt
(
"slideSize"
,
60
);
String
jdbcUrl
=
parameterTool
.
get
(
"jdbcUrl"
,
"jdbc:mysql://172.18.44.3:3306/jerry_test?user=root&password=5OqYM^zLwotJ3oSo&autoReconnect=true&useSSL=false"
);
Integer
maxRetry
=
parameterTool
.
getInt
(
"maxRetry"
,
3
);
Long
retryInteral
=
parameterTool
.
getLong
(
"retryInteral"
,
3000
);
String
checkpointPath
=
parameterTool
.
get
(
"checkpointPath"
,
"hdfs://bj-gmei-hdfs/user/data/flink/ctr-estimate/checkpoint"
);
Boolean
isStartFromEarliest
=
parameterTool
.
getBoolean
(
"isStartFromEarliest"
,
false
);
Boolean
isStartFromLatest
=
parameterTool
.
getBoolean
(
"isStartFromLatest"
,
false
);
String
startTime
=
parameterTool
.
get
(
"startTime"
);
Integer
parallelism
=
parameterTool
.
getInt
(
"parallelism"
,
2
);
String
zxJdbcUrl
=
parameterTool
.
get
(
"zxJdbcUrl"
,
"jdbc:mysql://172.16.30.141:3306/zhengxing?characterEncoding=UTF-8&autoReconnect=true&useSSL=false"
);
String
zxUsername
=
parameterTool
.
get
(
"zxUsername"
,
"work"
);
String
zxPassword
=
parameterTool
.
get
(
"zxPassword"
,
"BJQaT9VzDcuPBqkd"
);
String
jerryJdbcUrl
=
parameterTool
.
get
(
"jerryJdbcUrl"
,
"jdbc:mysql://172.16.40.170:4000/jerry_test?characterEncoding=UTF-8&autoReconnect=true&useSSL=false"
);
String
jerryUsername
=
parameterTool
.
get
(
"jerryUsername"
,
"data_user"
);
String
jerryPassword
=
parameterTool
.
get
(
"jerryPassword"
,
"YPEzp78HQBuhByWPpefQu6X3D6hEPfD6"
);
// 核心参数打印
System
.
out
.
println
(
"**********************************************************"
);
System
.
out
.
println
(
"*** inBrokers: "
+
inBrokers
);
System
.
out
.
println
(
"*** maidianInTopic: "
+
maidianInTopic
);
System
.
out
.
println
(
"*** maidianInGroupId: "
+
maidianInGroupId
);
System
.
out
.
println
(
"*** jdbcUrl: "
+
jdbcUrl
);
System
.
out
.
println
(
"*** checkpointPath: "
+
checkpointPath
);
System
.
out
.
println
(
"*** startTime: "
+
startTime
);
System
.
out
.
println
(
"*** windowSize: "
+
windowSize
);
System
.
out
.
println
(
"*** slideSize: "
+
slideSize
);
System
.
out
.
println
(
"*** zxJdbcUrl: "
+
zxJdbcUrl
);
System
.
out
.
println
(
"*** zxUsername: "
+
zxUsername
);
System
.
out
.
println
(
"*** zxPassword: "
+
zxPassword
);
System
.
out
.
println
(
"*** jerryJdbcUrl: "
+
jerryJdbcUrl
);
System
.
out
.
println
(
"*** jerryUsername: "
+
jerryUsername
);
System
.
out
.
println
(
"*** jerryPassword: "
+
jerryPassword
);
System
.
out
.
println
(
"**********************************************************"
);
// 获得流处理环境对象
StreamExecutionEnvironment
env
=
StreamExecutionEnvironment
.
getExecutionEnvironment
();
//env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
env
.
enableCheckpointing
(
1000
);
env
.
setStateBackend
(
new
FsStateBackend
(
checkpointPath
));
env
.
setRestartStrategy
(
RestartStrategies
.
fixedDelayRestart
(
1
,
3000
));
CheckpointConfig
config
=
env
.
getCheckpointConfig
();
config
.
enableExternalizedCheckpoints
(
CheckpointConfig
.
ExternalizedCheckpointCleanup
.
RETAIN_ON_CANCELLATION
);
DataStream
MaidianDataStream
=
new
MaidianKafkaSource
(
env
,
inBrokers
,
maidianInTopic
,
maidianInGroupId
,
batchSize
,
isStartFromEarliest
,
isStartFromLatest
,
startTime
).
getInstance
();
// 执行处理核心逻辑
new
CtrEstimateClkOperator
(
MaidianDataStream
,
jdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
,
windowSize
,
slideSize
).
run
();
new
CtrEstimateTagOperator
(
MaidianDataStream
,
jdbcUrl
,
maxRetry
,
retryInteral
,
parallelism
,
windowSize
,
slideSize
,
zxJdbcUrl
,
zxUsername
,
zxPassword
,
jerryJdbcUrl
,
jerryUsername
,
jerryPassword
).
run
();
// 常驻执行
env
.
execute
(
"ctr-estimate"
);
}
}
src/main/java/com/gmei/data/ctr/CtrEstimateMainClk.java
→
src/main/java/com/gmei/data/ctr/
Prod
CtrEstimateMainClk.java
View file @
3ba5e63c
...
...
@@ -10,13 +10,13 @@ import org.apache.flink.streaming.api.environment.CheckpointConfig;
import
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
;
/**
* @ClassName
CtrEstimateMain
* @ClassName
ProdCtrEstimateMainClk
* @Description: CTR预估特征实时处理入口
* @Author apple
* @Date 2020/3/30
* @Version V1.0
**/
public
class
CtrEstimateMainClk
{
public
class
Prod
CtrEstimateMainClk
{
public
static
void
main
(
String
[]
args
)
throws
Exception
{
// 获取运行参数
...
...
src/main/java/com/gmei/data/ctr/CtrEstimateMainTag.java
→
src/main/java/com/gmei/data/ctr/
Prod
CtrEstimateMainTag.java
View file @
3ba5e63c
...
...
@@ -10,13 +10,13 @@ import org.apache.flink.streaming.api.environment.CheckpointConfig;
import
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
;
/**
* @ClassName
CtrEstimateMain
* @ClassName
ProdCtrEstimateMainTag
* @Description: CTR预估特征实时处理入口
* @Author apple
* @Date 2020/3/30
* @Version V1.0
**/
public
class
CtrEstimateMainTag
{
public
class
Prod
CtrEstimateMainTag
{
public
static
void
main
(
String
[]
args
)
throws
Exception
{
// 获取运行参数
...
...
src/main/java/com/gmei/data/ctr/
CtrEstimateMainClkDev
.java
→
src/main/java/com/gmei/data/ctr/
TestCtrEstimateMainClk
.java
View file @
3ba5e63c
...
...
@@ -10,13 +10,13 @@ import org.apache.flink.streaming.api.environment.CheckpointConfig;
import
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
;
/**
* @ClassName
CtrEstimateMain
* @ClassName
TestCtrEstimateMainClk
* @Description: CTR预估特征实时处理入口
* @Author apple
* @Date 2020/3/30
* @Version V1.0
**/
public
class
CtrEstimateMainClkDev
{
public
class
TestCtrEstimateMainClk
{
public
static
void
main
(
String
[]
args
)
throws
Exception
{
// 获取运行参数
...
...
src/main/java/com/gmei/data/ctr/
CtrEstimateMainTagDev
.java
→
src/main/java/com/gmei/data/ctr/
TestCtrEstimateMainTag
.java
View file @
3ba5e63c
...
...
@@ -7,13 +7,13 @@ import org.apache.flink.streaming.api.datastream.DataStream;
import
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
;
/**
* @ClassName
CtrEstimateMain
* @ClassName
TestCtrEstimateMainTag
* @Description: CTR预估特征实时处理入口
* @Author apple
* @Date 2020/3/30
* @Version V1.0
**/
public
class
CtrEstimateMainTagDev
{
public
class
TestCtrEstimateMainTag
{
public
static
void
main
(
String
[]
args
)
throws
Exception
{
// 获取运行参数
...
...
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