Commit 5f9a82bd authored by sheteng's avatar sheteng

调整数据 offset

parent bb95876a
...@@ -57,10 +57,10 @@ dependencies { ...@@ -57,10 +57,10 @@ dependencies {
// https://mvnrepository.com/artifact/mysql/mysql-connector-java // https://mvnrepository.com/artifact/mysql/mysql-connector-java
implementation group: 'mysql', name: 'mysql-connector-java', version: '8.0.25' implementation group: 'mysql', name: 'mysql-connector-java', version: '8.0.25'
compileOnly "org.apache.flink:flink-statebackend-rocksdb_${scalaBinaryVersion}:${flinkVersion}" compile "org.apache.flink:flink-statebackend-rocksdb_${scalaBinaryVersion}:${flinkVersion}"
compileOnly "org.apache.flink:flink-java:${flinkVersion}" compile "org.apache.flink:flink-java:${flinkVersion}"
compileOnly "org.apache.flink:flink-streaming-java_${scalaBinaryVersion}:${flinkVersion}" compile "org.apache.flink:flink-streaming-java_${scalaBinaryVersion}:${flinkVersion}"
compileOnly "org.apache.flink:flink-clients_${scalaBinaryVersion}:${flinkVersion}" compile "org.apache.flink:flink-clients_${scalaBinaryVersion}:${flinkVersion}"
implementation 'ru.yandex.clickhouse:clickhouse-jdbc:0.1.52' implementation 'ru.yandex.clickhouse:clickhouse-jdbc:0.1.52'
// -------------------------------------------------------------- // --------------------------------------------------------------
// Dependencies that should be part of the shadow jar, e.g. // Dependencies that should be part of the shadow jar, e.g.
......
...@@ -31,6 +31,7 @@ public class AppsflyerAnalyze { ...@@ -31,6 +31,7 @@ public class AppsflyerAnalyze {
appConfig.setSourceKafkaBootstrapServers(parameters.get("source-kafka-bootstrap-servers")); appConfig.setSourceKafkaBootstrapServers(parameters.get("source-kafka-bootstrap-servers"));
appConfig.setSourceKafkaGroupId(parameters.get("source-kafka-group-id")); appConfig.setSourceKafkaGroupId(parameters.get("source-kafka-group-id"));
appConfig.setSourceKafkaTopic(parameters.get("source-kafka-topic")); appConfig.setSourceKafkaTopic(parameters.get("source-kafka-topic"));
appConfig.setStartOffsetTime(Long.valueOf(parameters.get("start-offset-time")));
appConfig.setSinkKafkaBootstrapServers(parameters.get("sink-kafka-bootstrap-servers")); appConfig.setSinkKafkaBootstrapServers(parameters.get("sink-kafka-bootstrap-servers"));
appConfig.setSinkKafkaTopic(parameters.get("sink-kafka-topic")); appConfig.setSinkKafkaTopic(parameters.get("sink-kafka-topic"));
...@@ -68,8 +69,10 @@ public class AppsflyerAnalyze { ...@@ -68,8 +69,10 @@ public class AppsflyerAnalyze {
properties.setProperty("group.id", appConfig.getSourceKafkaGroupId()); properties.setProperty("group.id", appConfig.getSourceKafkaGroupId());
properties.setProperty("auto.offset.reset", "earliest"); properties.setProperty("auto.offset.reset", "earliest");
FlinkKafkaConsumer<String> dataConsumer = new FlinkKafkaConsumer<>(appConfig.getSourceKafkaTopic(), new SimpleStringSchema(), properties); FlinkKafkaConsumer<String> dataConsumer = new FlinkKafkaConsumer<>(appConfig.getSourceKafkaTopic(), new SimpleStringSchema(), properties);
// TODO: 2021/8/10 设置消费开始的时间戳1628524944000 // TODO: 2021/8/10 设置消费开始的时间戳1632326555000
// dataConsumer.setStartFromTimestamp(1631548945000L); if (appConfig.getStartOffsetTime() > 0) {
dataConsumer.setStartFromTimestamp(appConfig.getStartOffsetTime());
}
Properties props2 = new Properties(); Properties props2 = new Properties();
props2.setProperty("bootstrap.servers", appConfig.getSinkKafkaBootstrapServers()); props2.setProperty("bootstrap.servers", appConfig.getSinkKafkaBootstrapServers());
......
...@@ -46,4 +46,6 @@ public class AppConfig implements Serializable { ...@@ -46,4 +46,6 @@ public class AppConfig implements Serializable {
private String sinkType; private String sinkType;
private Long startOffsetTime;
} }
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment