通过SpringBoot联合Flink的高性能流处理本事,开发者可以构建出高效、可靠的及时数据处理办理方案。这些应用不仅可以或许提高企业的运营服从,还能增强竞争力。
应用场景
1. 推荐系统:
根据用户的及时行为生成个性化推荐。
更新用户画像,以适应不断变化的兴趣和偏好。
2. 变乱驱动架构:
3. 欺诈检测:
4. 流式数据分析:
及时分析传感器数据,如物联网设备的数据流。
监控和分析用户行为数据,提供即时反馈。
5. 日志处理:
6. 及时报表:
生成及时销售陈诉、市场趋势分析等。
支持决议订定者获得最新的业务洞察力。
7. 供应链管理:
及时跟踪库存水平、订单状态和物流信息。
主动调整生产计划以满足需求变化。
8. 交际媒体分析:
及时监测舆情,分析公众对特定话题的见解。
提供及时的情感分析,帮助企业相识品牌形象。
9. 网络安全:
及时监控网络流量,检测埋伏的安全威胁。
快速相应攻击,保护企业免受损害。
代码演示
我们现在用简朴代码例子演示怎样根据用户的及时行为生成个性化推荐。以下代码将展示怎样从Kafka读取用户行为数据,并使用Flink进行及时处理,最后生成个性化的推荐效果并将其写回到另一个Kafka主题中。
1. 依赖项
在pom.xml中添加须要的依赖项:
- <dependencies>
- <!-- Spring Boot Starter Web -->
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-web</artifactId>
- </dependency>
- <!-- Flink dependencies -->
- <dependency>
- <groupId>org.apache.flink</groupId>
- <artifactId>flink-java</artifactId>
- <version>1.14.6</version>
- </dependency>
- <dependency>
- <groupId>org.apache.flink</groupId>
- <artifactId>flink-streaming-java_2.12</artifactId>
- <version>1.14.6</version>
- </dependency>
- <dependency>
- <groupId>org.apache.flink</groupId>
- <artifactId>flink-clients_2.12</artifactId>
- <version>1.14.6</version>
- </dependency>
- <dependency>
- <groupId>org.apache.flink</groupId>
- <artifactId>flink-connector-kafka_2.12</artifactId>
- <version>1.14.6</version>
- </dependency>
- <!-- JSON processing with Jackson -->
- <dependency>
- <groupId>com.fasterxml.jackson.core</groupId>
- <artifactId>jackson-databind</artifactId>
- </dependency>
- </dependencies>
复制代码 2. 界说数据模型
界说用户行为数据和推荐效果的数据模型。
- package com.example.demo.model;
- publicclass UserBehavior {
- private String userId;
- private String productId;
- private String action; // e.g., "view", "click", "purchase"
- privatelong timestamp;
- // Getters and setters
- public String getUserId() {
- return userId;
- }
- public void setUserId(String userId) {
- this.userId = userId;
- }
- public String getProductId() {
- return productId;
- }
- public void setProductId(String productId) {
- this.productId = productId;
- }
- public String getAction() {
- return action;
- }
- public void setAction(String action) {
- this.action = action;
- }
- public long getTimestamp() {
- return timestamp;
- }
- public void setTimestamp(long timestamp) {
- this.timestamp = timestamp;
- }
- }
复制代码- package com.example.demo.model;
- import java.util.List;
- publicclass RecommendationResult {
- private String userId;
- private List<String> recommendedProducts;
- // Getters and setters
- public String getUserId() {
- return userId;
- }
- public void setUserId(String userId) {
- this.userId = userId;
- }
- public List<String> getRecommendedProducts() {
- return recommendedProducts;
- }
- public void setRecommendedProducts(List<String> recommendedProducts) {
- this.recommendedProducts = recommendedProducts;
- }
- }
复制代码 3. 编写Flink Job
创建一个简朴的Flink作业,它从Kafka读取用户行为数据,计算每个用户的热门产品,并将推荐效果写入另一个Kafka主题。
- package com.example.demo;
- import com.example.demo.model.RecommendationResult;
- import com.example.demo.model.UserBehavior;
- import org.apache.flink.api.common.functions.MapFunction;
- import org.apache.flink.api.java.tuple.Tuple2;
- import org.apache.flink.streaming.api.datastream.DataStream;
- import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
- import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
- import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
- import org.apache.flink.api.common.serialization.SimpleStringSchema;
- import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;
- import java.util.ArrayList;
- import java.util.HashMap;
- import java.util.List;
- import java.util.Map;
- import java.util.Properties;
- publicclass RecommendationJob {
- public static void main(String[] args) throws Exception {
- final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
- Properties properties = new Properties();
- properties.setProperty("bootstrap.servers", "localhost:9092");
- properties.setProperty("group.id", "test");
- FlinkKafkaConsumer<String> kafkaConsumer = new FlinkKafkaConsumer<>("user-behavior-topic", new SimpleStringSchema(), properties);
- kafkaConsumer.setStartFromEarliest();
- ObjectMapper objectMapper = new ObjectMapper();
- DataStream<UserBehavior> userBehaviors = env.addSource(kafkaConsumer)
- .map((MapFunction<String, UserBehavior>) value -> objectMapper.readValue(value, UserBehavior.class));
- DataStream<Tuple2<String, Map<String, Integer>>> productCountsPerUser = userBehaviors
- .filter(behavior -> behavior.getAction().equals("purchase"))
- .keyBy(UserBehavior::getUserId)
- .flatMap((FlatMapFunction<UserBehavior, Tuple2<String, Map<String, Integer>>>) (value, out) -> {
- Map<String, Integer> productCounts = new HashMap<>();
- if (productCounts.containsKey(value.getProductId())) {
- productCounts.put(value.getProductId(), productCounts.get(value.getProductId()) + 1);
- } else {
- productCounts.put(value.getProductId(), 1);
- }
- out.collect(Tuple2.of(value.getUserId(), productCounts));
- })
- .keyBy(Tuple2::f0)
- .reduce((t1, t2) -> {
- t1.f1.forEach((productId, count) -> t2.f1.merge(productId, count, Integer::sum));
- return t1;
- });
- DataStream<RecommendationResult> recommendations = productCountsPerUser
- .map((MapFunction<Tuple2<String, Map<String, Integer>>, RecommendationResult>) value -> {
- RecommendationResult result = new RecommendationResult();
- result.setUserId(value.f0);
- List<String> recommendedProducts = new ArrayList<>(value.f1.keySet());
- result.setRecommendedProducts(recommendedProducts.subList(0, Math.min(5, recommendedProducts.size())));
- return result;
- });
- FlinkKafkaProducer<String> kafkaProducer = new FlinkKafkaProducer<>("recommendations-topic",
- (RecommendationResult recommendation) -> objectMapper.writeValueAsString(recommendation),
- properties);
- recommendations.map(ObjectMapper::writeValueAsString).addSink(kafkaProducer);
- env.execute("Recommendation Job");
- }
- }
复制代码 4. 创建Spring Boot Application
创建一个简朴的Spring Boot应用步伐来启动Flink作业。
- package com.example.demo;
- import org.springframework.boot.SpringApplication;
- import org.springframework.boot.autoconfigure.SpringBootApplication;
- import org.springframework.context.annotation.Bean;
- @SpringBootApplication
- publicclass DemoApplication {
- public static void main(String[] args) {
- SpringApplication.run(DemoApplication.class, args);
- }
- @Bean
- public Runnable flinkRunner() {
- return () -> {
- try {
- RecommendationJob.main(new String[]{});
- } catch (Exception e) {
- thrownew RuntimeException(e);
- }
- };
- }
- }
复制代码 测试效果
“ 请确保你的环境中已经设置好了Kafka和其他须要的服务,以便代码可以或许正常运行。
我本地的Kafka集群正在运行,并且有两个主题:user-behavior-topic 和 recommendations-topic。
关注我,送Java福利
- /**
- * 这段代码只有Java开发者才能看得懂!
- * 关注我微信公众号之后,
- * 发送:"666",
- * 即可获得一本由Java大神一手面试经验诚意出品
- * 《Java开发者面试百宝书》Pdf电子书
- * 福利截止日期为2025年01月22日止
- * 手快有手慢没!!!
- */
- System.out.println("请关注我的微信公众号:");
- System.out.println("Java知识日历");
复制代码 免责声明:如果侵犯了您的权益,请联系站长,我们会及时删除侵权内容,谢谢合作!更多信息从访问主页:qidao123.com:ToB企服之家,中国第一个企服评测及商务社交产业平台。 |