Commit b5ec291f by lichaomin

搭建

parent 69645267
...@@ -13,5 +13,55 @@ ...@@ -13,5 +13,55 @@
<groupId>com.byit.byit-screen</groupId> <groupId>com.byit.byit-screen</groupId>
<artifactId>byit-screen-computer</artifactId> <artifactId>byit-screen-computer</artifactId>
<properties>
<flink.version>1.8.0</flink.version>
<scala.binary.version>2.11</scala.binary.version>
<byit.common.util>0.0.1-SNAPSHOT</byit.common.util>
<byit-common-kafka>0.0.1-SNAPSHOT</byit-common-kafka>
</properties>
<dependencies>
<dependency>
<groupId>com.byit.common</groupId>
<artifactId>byit-common-utils</artifactId>
<version>${byit.common.util}</version>
</dependency>
<dependency>
<groupId>com.byit.common</groupId>
<artifactId>byit-common-kafka</artifactId>
<version>${byit-common-kafka}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka-0.11_${scala.binary.version}</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java_${scala.binary.version}</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java_2.11</artifactId>
<version>1.8.0</version>
</dependency>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework</groupId>
<artifactId>spring-test</artifactId>
<version>4.3.14.RELEASE</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
</project> </project>
\ No newline at end of file
package com.byit.computer.constant;
/**
* lcm
*/
public class PropertiesConstants {
public static final String KAFKA_BROKERS = "kafka.brokers";
public static final String KAFKA_ZOOKEEPER_CONNECT = "kafka.zookeeper.connect";
public static final String KAFKA_GROUP_ID = "kafka.group.id";
public static final String METRICS_TOPIC = "kafka.producer.topic";
public static final String CONSUMER_FROM_TIME = "consumer.from.time";
public static final String STREAM_PARALLELISM = "stream.parallelism";
public static final String STREAM_SINK_PARALLELISM = "stream.sink.parallelism";
public static final String STREAM_DEFAULT_PARALLELISM = "stream.default.parallelism";
public static final String STREAM_CHECKPOINT_ENABLE = "stream.checkpoint.enable";
public static final String STREAM_CHECKPOINT_INTERVAL = "stream.checkpoint.interval";
}
package com.byit.computer.resource;
import com.byit.computer.constant.PropertiesConstants;
import org.apache.flink.api.common.restartstrategy.RestartStrategies;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.TimeCharacteristic;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer011;
import org.apache.flink.streaming.connectors.kafka.internals.KafkaTopicPartition;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndTimestamp;
import org.apache.kafka.common.PartitionInfo;
import org.apache.kafka.common.TopicPartition;
import org.springframework.beans.factory.annotation.Value;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;
/**
* lcm
* flink input
*/
public class KafkaSource {
/**
* 设置基础的 Kafka 配置
*
* @return
*/
public static Properties buildKafkaProps() {
return buildKafkaProps(ParameterTool.fromSystemProperties());
}
/**
* 设置 kafka 配置
*
* @param parameterTool
* @return
*/
public static Properties buildKafkaProps(ParameterTool parameterTool) {
Properties props = parameterTool.getProperties();
props.put("bootstrap.servers", parameterTool.get(PropertiesConstants.KAFKA_BROKERS, "10.0.10.136:9092,10.0.10.137:9092,10.0.10.138:9092"));
props.put("zookeeper.connect", parameterTool.get(PropertiesConstants.KAFKA_ZOOKEEPER_CONNECT, "10.0.10.136:2181,10.0.10.137:2181,10.0.10.138:2181"));
props.put("group.id", parameterTool.get(PropertiesConstants.KAFKA_GROUP_ID, "test"));
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("auto.offset.reset", "latest");
return props;
}
public static DataStreamSource<String> buildSource(StreamExecutionEnvironment env) throws IllegalAccessException {
ParameterTool parameter = (ParameterTool) env.getConfig().getGlobalJobParameters();
String topic = parameter.getRequired(PropertiesConstants.METRICS_TOPIC);
Long time = parameter.getLong(PropertiesConstants.CONSUMER_FROM_TIME, 1L);
return buildSource(env, topic, time);
}
/**
* @param env
* @param topic
* @param time 订阅的时间
* @return
* @throws IllegalAccessException
*/
public static DataStreamSource<String> buildSource(StreamExecutionEnvironment env, String topic, Long time) throws IllegalAccessException {
ParameterTool parameterTool = (ParameterTool) env.getConfig().getGlobalJobParameters();
Properties props = buildKafkaProps(parameterTool);
FlinkKafkaConsumer011<String> consumer = new FlinkKafkaConsumer011<>(
topic,
new SimpleStringSchema(),
props);
//重置offset到time时刻
if (time != 0L) {
Map<KafkaTopicPartition, Long> partitionOffset = buildOffsetByTime(props, parameterTool, time);
consumer.setStartFromSpecificOffsets(partitionOffset);
}
return env.addSource(consumer);
}
private static Map<KafkaTopicPartition, Long> buildOffsetByTime(Properties props, ParameterTool parameterTool, Long time) {
props.setProperty("group.id", "query_time_" + time);
KafkaConsumer consumer = new KafkaConsumer(props);
List<PartitionInfo> partitionsFor = consumer.partitionsFor(parameterTool.getRequired(PropertiesConstants.METRICS_TOPIC));
Map<TopicPartition, Long> partitionInfoLongMap = new HashMap<>();
for (PartitionInfo partitionInfo : partitionsFor) {
partitionInfoLongMap.put(new TopicPartition(partitionInfo.topic(), partitionInfo.partition()), time);
}
Map<TopicPartition, OffsetAndTimestamp> offsetResult = consumer.offsetsForTimes(partitionInfoLongMap);
Map<KafkaTopicPartition, Long> partitionOffset = new HashMap<>();
offsetResult.forEach((key, value) -> partitionOffset.put(new KafkaTopicPartition(key.topic(), key.partition()), value.offset()));
consumer.close();
return partitionOffset;
}
}
package com.byit.computer.sink;
import com.byit.computer.resource.KafkaSource;
import com.byit.computer.util.ExecutionEnvUtil;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer011;
import org.apache.flink.streaming.connectors.kafka.internals.KeyedSerializationSchemaWrapper;
/**
* lcm
* out
*/
public class KafkaSink {
public static void setKafkaSink(DataStreamSource data,ParameterTool parameterTool) throws Exception {
data.addSink(new FlinkKafkaProducer011<>(
parameterTool.get("spring.kafka.bootstrap-servers"),
parameterTool.get("kafka.consumer.topic"),
new KeyedSerializationSchemaWrapper<>(new SimpleStringSchema())
)).name("flink-connectors-kafka")
.setParallelism(parameterTool.getInt("stream.sink.parallelism"));
}
}
package com.byit.computer.util;
import com.byit.common.utils.Constants;
import com.byit.computer.constant.PropertiesConstants;
import org.apache.flink.api.common.restartstrategy.RestartStrategies;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.TimeCharacteristic;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
/**
* lcm
*/
public class ExecutionEnvUtil {
public static ParameterTool createParameterTool(final String[] args) throws Exception {
return ParameterTool
.fromPropertiesFile(ExecutionEnvUtil.class.getResourceAsStream(Constants.CHAR_SLASH+Constants.CONFIG+"/"+System.getenv(Constants.ENV_HOME)+"/"+Constants.APPLICATION))
.mergeWith(ParameterTool.fromArgs(args))
.mergeWith(ParameterTool.fromSystemProperties())
.mergeWith(ParameterTool.fromMap(getenv()));
}
public static void prepare(StreamExecutionEnvironment env,ParameterTool parameterTool) throws Exception {
env.setParallelism(parameterTool.getInt(PropertiesConstants.STREAM_PARALLELISM, 5));
env.getConfig().disableSysoutLogging();
env.getConfig().setRestartStrategy(RestartStrategies.fixedDelayRestart(4, 10000));
if (parameterTool.getBoolean(PropertiesConstants.STREAM_CHECKPOINT_ENABLE, true)) {
env.enableCheckpointing(parameterTool.getInt(PropertiesConstants.STREAM_CHECKPOINT_INTERVAL, 1000)); // create a checkpoint every 5 seconds
}
env.getConfig().setGlobalJobParameters(parameterTool); // make parameters available in the web interface
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
}
private static Map<String, String> getenv() {
Map<String, String> map = new HashMap<>();
for (Map.Entry<String, String> entry : System.getenv().entrySet()) {
map.put(entry.getKey(), entry.getValue());
}
return map;
}
}
...@@ -23,8 +23,24 @@ logging.config=classpath:config/logging-config.xml ...@@ -23,8 +23,24 @@ logging.config=classpath:config/logging-config.xml
logging.path=/home/logger/screen_service_log logging.path=/home/logger/screen_service_log
logging.file=screen_service_log logging.file=screen_service_log
spring.kafka.bootstrap-servers=10.0.10.136:9092,10.0.10.137:9092,10.0.10.138:9092
kafka.producer.topic=input
kafka.producer.retries=0
kafka.producer.batch.size=4096
kafka.producer.linger=1
kafka.producer.buffer.memory=40960
kafka.consumer.zookeeper.connect=10.0.10.136:2181,10.0.10.137:2181,10.0.10.138:2181
kafka.consumer.servers=10.0.10.136:9092,10.0.10.137:9092,10.0.10.138:9092
kafka.consumer.enable.auto.commit=true
kafka.consumer.session.timeout=6000
kafka.consumer.auto.commit.interval=100
kafka.consumer.auto.offset.reset=latest
kafka.consumer.topic=output
kafka.consumer.group.id=output
kafka.consumer.concurrency=10
stream.sink.parallelism=4
\ No newline at end of file
package com.byit;
import com.byit.computer.resource.KafkaSource;
import com.byit.computer.sink.KafkaSink;
import com.byit.computer.util.ExecutionEnvUtil;
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.common.functions.ReduceFunction;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer011;
import org.apache.flink.streaming.connectors.kafka.internals.KeyedSerializationSchemaWrapper;
import org.apache.flink.util.Collector;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.test.context.junit4.SpringRunner;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
public class TestApplication {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
final ParameterTool parameterTool = ExecutionEnvUtil.createParameterTool(args);
ExecutionEnvUtil.prepare(env,parameterTool);
DataStreamSource<String> data = KafkaSource.buildSource(env);
data.flatMap(new FlatMapFunction<String, Tuple2<String, Integer>>() {
@Override
public void flatMap(String value, Collector<Tuple2<String, Integer>> out) throws Exception {
String[] splits = value.toLowerCase().split("\\W+");
for (String split : splits) {
if (split.length() > 0) {
out.collect(new Tuple2<>(split, 1));
}
}
}
}).groupBy(0)
.reduce(new ReduceFunction<Tuple2<String, Integer>>() {
@Override
public Tuple2<String, Integer> reduce(Tuple2<String, Integer> value1, Tuple2<String, Integer> value2) throws Exception {
return new Tuple2<>(value1.f0, value1.f1 + value1.f1);
}
})
.print();
KafkaSink.setKafkaSink(data,parameterTool);
//data.print();
env.execute("flink learning connectors kafka");
}
}
package com.byit;
import org.apache.kafka.clients.producer.KafkaProducer;
import java.util.Properties;
public class TestProducer {
public static final String broker_list = "10.0.10.136:9092,10.0.10.137:9092,10.0.10.138:9092";
public static final String topic = "student"; //kafka topic 需要和 flink 程序用同一个 topic
public static void writeToKafka() throws InterruptedException {
Properties props = new Properties();
props.put("bootstrap.servers", broker_list);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
KafkaProducer producer = new KafkaProducer<String, String>(props);
for (int i = 1; i <= 10000; i++) {
/*Student student = new Student(i, "zhisheng" + i, "password" + i, 18 + i);
ProducerRecord record = new ProducerRecord<String, String>(topic, null, null, GsonUtil.toJson(student));
producer.send(record);
System.out.println("发送数据: " + GsonUtil.toJson(student));*/
}
producer.flush();
}
public static void main(String[] args) throws InterruptedException {
writeToKafka();
}
}
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