Java实现埋点数据收集案例

wen java案例 1

本文目录导读:

Java实现埋点数据收集案例

  1. 项目结构
  2. Maven依赖配置 (pom.xml)
  3. 数据模型
  4. 控制器
  5. 服务层
  6. Repository
  7. 工具类
  8. 配置文件
  9. 前端埋点代码示例
  10. 数据库表结构
  11. 数据统计接口示例

我将为您提供一个完整的Java埋点数据收集系统案例,包含服务端接收、数据处理和存储的完整实现。

项目结构

analytics-system/
├── pom.xml
├── src/main/java/com/analytics/
│   ├── controller/
│   │   └── TrackController.java
│   ├── service/
│   │   ├── TrackService.java
│   │   ├── TrackServiceImpl.java
│   │   └── KafkaProducerService.java
│   ├── model/
│   │   ├── TrackEvent.java
│   │   └── UserAction.java
│   ├── config/
│   │   ├── KafkaConfig.java
│   │   └── WebConfig.java
│   ├── repository/
│   │   └── TrackEventRepository.java
│   └── utils/
│       ├── IpUtils.java
│       └── UserAgentUtils.java

Maven依赖配置 (pom.xml)

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 
         http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <groupId>com.analytics</groupId>
    <artifactId>analytics-system</artifactId>
    <version>1.0.0</version>
    <packaging>jar</packaging>
    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>2.7.14</version>
    </parent>
    <properties>
        <java.version>1.8</java.version>
        <kafka.version>2.8.0</kafka.version>
    </properties>
    <dependencies>
        <!-- Spring Boot Web -->
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>
        <!-- Spring Data JPA -->
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-data-jpa</artifactId>
        </dependency>
        <!-- MySQL -->
        <dependency>
            <groupId>mysql</groupId>
            <artifactId>mysql-connector-java</artifactId>
            <version>8.0.33</version>
        </dependency>
        <!-- Kafka -->
        <dependency>
            <groupId>org.springframework.kafka</groupId>
            <artifactId>spring-kafka</artifactId>
        </dependency>
        <!-- Redis -->
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-data-redis</artifactId>
        </dependency>
        <!-- Lombok -->
        <dependency>
            <groupId>org.projectlombok</groupId>
            <artifactId>lombok</artifactId>
            <optional>true</optional>
        </dependency>
        <!-- User-Agent解析 -->
        <dependency>
            <groupId>eu.bitwalker</groupId>
            <artifactId>UserAgentUtils</artifactId>
            <version>1.21</version>
        </dependency>
        <!-- Jackson JSON处理 -->
        <dependency>
            <groupId>com.fasterxml.jackson.core</groupId>
            <artifactId>jackson-databind</artifactId>
        </dependency>
    </dependencies>
</project>

数据模型

TrackEvent.java - 埋点事件模型

package com.analytics.model;
import lombok.Data;
import javax.persistence.*;
import java.time.LocalDateTime;
@Data
@Entity
@Table(name = "track_events")
public class TrackEvent {
    @Id
    @GeneratedValue(strategy = GenerationType.IDENTITY)
    private Long id;
    // 事件ID(前端生成的唯一标识)
    @Column(name = "event_id", unique = true)
    private String eventId;
    // 事件类型:click, view, page_view, custom等
    @Column(name = "event_type", nullable = false)
    private String eventType;
    // 事件名称
    @Column(name = "event_name")
    private String eventName;
    // 用户ID
    @Column(name = "user_id")
    private String userId;
    // 用户设备ID
    @Column(name = "device_id")
    private String deviceId;
    // 页面URL
    @Column(name = "page_url")
    private String pageUrl;
    // 页面标题
    @Column(name = "page_title")
    private String pageTitle;
    // 元素ID(被点击的元素)
    @Column(name = "element_id")
    private String elementId;
    // 元素类名
    @Column(name = "element_class")
    private String elementClass;
    // 元素文本内容
    @Column(name = "element_text")
    private String elementText;
    // 业务数据JSON
    @Column(name = "business_data", columnDefinition = "TEXT")
    private String businessData;
    // IP地址
    @Column(name = "ip")
    private String ip;
    // 用户代理
    @Column(name = "user_agent")
    private String userAgent;
    // 设备类型:mobile, tablet, desktop
    @Column(name = "device_type")
    private String deviceType;
    // 操作系统
    @Column(name = "os")
    private String os;
    // 浏览器
    @Column(name = "browser")
    private String browser;
    // 浏览器版本
    @Column(name = "browser_version")
    private String browserVersion;
    // 访问来源
    @Column(name = "referrer")
    private String referrer;
    // 会话ID
    @Column(name = "session_id")
    private String sessionId;
    // 访问时间戳(毫秒)
    @Column(name = "timestamp")
    private Long timestamp;
    // 服务器接收时间
    @Column(name = "server_time")
    private LocalDateTime serverTime;
    // 客户端时间
    @Column(name = "client_time")
    private LocalDateTime clientTime;
    // 数据入库时间
    @Column(name = "created_at")
    private LocalDateTime createdAt;
    @PrePersist
    public void prePersist() {
        createdAt = LocalDateTime.now();
    }
}

UserAction.java - 用户行为请求

package com.analytics.model;
import lombok.Data;
import java.util.Map;
@Data
public class UserAction {
    // 事件类型
    private String eventType;
    // 事件名称
    private String eventName;
    // 事件数据
    private Map<String, Object> eventData;
    // 页面信息
    private String pageUrl;
    private String pageTitle;
    private String referrer;
    // 元素信息
    private String elementId;
    private String elementClass;
    private String elementText;
    // 时间戳
    private Long timestamp;
    // 会话ID
    private String sessionId;
    // 用户ID
    private String userId;
}

控制器

TrackController.java

package com.analytics.controller;
import com.analytics.model.TrackEvent;
import com.analytics.model.UserAction;
import com.analytics.service.TrackService;
import com.analytics.utils.IpUtils;
import com.analytics.utils.UserAgentUtils;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import javax.servlet.http.HttpServletRequest;
import java.util.UUID;
@RestController
@RequestMapping("/api/track")
public class TrackController {
    @Autowired
    private TrackService trackService;
    @Autowired
    private ObjectMapper objectMapper;
    /**
     * 接收埋点数据
     * 支持JSON格式和表单格式
     */
    @PostMapping("/collect")
    public Result collect(@RequestBody(required = false) UserAction userAction,
                          @RequestParam(required = false) String data,
                          HttpServletRequest request) {
        try {
            // 如果是字符串数据,解析JSON
            if (userAction == null && data != null) {
                userAction = objectMapper.readValue(data, UserAction.class);
            }
            if (userAction == null) {
                return Result.error("无效的埋点数据");
            }
            // 转换为TrackEvent
            TrackEvent event = convertToTrackEvent(userAction, request);
            // 异步处理埋点数据
            trackService.processTrackEvent(event);
            // 返回成功,包含事件ID
            return Result.success("埋点数据接收成功", event.getEventId());
        } catch (Exception e) {
            return Result.error("埋点数据接收失败: " + e.getMessage());
        }
    }
    /**
     * GET请求支持(用于简单埋点)
     */
    @GetMapping("/collect")
    public Result collectGet(@RequestParam String eventType,
                            @RequestParam String eventName,
                            HttpServletRequest request) {
        UserAction action = new UserAction();
        action.setEventType(eventType);
        action.setEventName(eventName);
        action.setTimestamp(System.currentTimeMillis());
        return collect(action, null, request);
    }
    /**
     * 批量接收埋点数据
     */
    @PostMapping("/collect/batch")
    public Result collectBatch(@RequestBody List<UserAction> actions,
                              HttpServletRequest request) {
        for (UserAction action : actions) {
            collect(action, null, request);
        }
        return Result.success("批量埋点数据接收成功");
    }
    /**
     * 转换处理
     */
    private TrackEvent convertToTrackEvent(UserAction userAction, HttpServletRequest request) {
        TrackEvent event = new TrackEvent();
        // 基本事件信息
        event.setEventId(UUID.randomUUID().toString());
        event.setEventType(userAction.getEventType());
        event.setEventName(userAction.getEventName());
        event.setPageUrl(userAction.getPageUrl());
        event.setPageTitle(userAction.getPageTitle());
        event.setReferrer(userAction.getReferrer());
        event.setElementId(userAction.getElementId());
        event.setElementClass(userAction.getElementClass());
        event.setElementText(userAction.getElementText());
        event.setUserId(userAction.getUserId());
        event.setSessionId(userAction.getSessionId());
        // 时间戳
        event.setTimestamp(userAction.getTimestamp() != null ? 
                          userAction.getTimestamp() : System.currentTimeMillis());
        // 设置客户端时间
        event.setClientTime(LocalDateTime.now());
        // 业务数据转换为JSON
        try {
            if (userAction.getEventData() != null) {
                event.setBusinessData(objectMapper.writeValueAsString(userAction.getEventData()));
            }
        } catch (Exception e) {
            event.setBusinessData("{}");
        }
        // 获取客户端信息
        event.setIp(IpUtils.getIpAddr(request));
        String userAgent = request.getHeader("User-Agent");
        event.setUserAgent(userAgent);
        // 解析User-Agent
        UserAgentUtils.parseUserAgent(userAgent, event);
        return event;
    }
}
// 统一的响应结果
class Result {
    private int code;
    private String message;
    private Object data;
    public static Result success(String message) {
        return success(message, null);
    }
    public static Result success(String message, Object data) {
        Result result = new Result();
        result.code = 200;
        result.message = message;
        result.data = data;
        return result;
    }
    public static Result error(String message) {
        Result result = new Result();
        result.code = 500;
        result.message = message;
        return result;
    }
    // getters and setters...
}

服务层

TrackService.java

package com.analytics.service;
import com.analytics.model.TrackEvent;
public interface TrackService {
    /**
     * 处理埋点事件
     */
    void processTrackEvent(TrackEvent event);
    /**
     * 统计数据
     */
    EventStats getEventStats(String eventType);
    /**
     * 按时间段统计
     */
    List<EventStats> getEventStatsByTimeRange(LocalDateTime start, LocalDateTime end);
    // ... 其他统计方法
}

TrackServiceImpl.java

package com.analytics.service;
import com.analytics.model.TrackEvent;
import com.analytics.repository.TrackEventRepository;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@Service
public class TrackServiceImpl implements TrackService {
    private static final Logger logger = LoggerFactory.getLogger(TrackServiceImpl.class);
    @Autowired
    private TrackEventRepository trackEventRepository;
    @Autowired
    private KafkaProducerService kafkaProducerService;
    /**
     * 异步处理埋点数据
     */
    @Override
    @Async("asyncExecutor")
    @Transactional
    public void processTrackEvent(TrackEvent event) {
        try {
            logger.info("处理埋点数据: {}", event.getEventId());
            // 1. 发送到Kafka进行实时处理
            kafkaProducerService.sendTrackEvent(event);
            // 2. 保存到数据库
            trackEventRepository.save(event);
            // 3. 缓存到Redis(用于实时统计)
            cacheToRedis(event);
        } catch (Exception e) {
            logger.error("处理埋点数据失败: {}", event.getEventId(), e);
            // 可以发送到失败队列或做其他处理
        }
    }
    /**
     * 缓存到Redis
     */
    private void cacheToRedis(TrackEvent event) {
        try {
            String key = "track:" + event.getEventType() + ":" + 
                         DateUtils.formatDate(event.getTimestamp(), "yyyyMMdd");
            redisTemplate.opsForList().leftPush(key, event);
            // 设置过期时间(例如7天)
            redisTemplate.expire(key, 7, TimeUnit.DAYS);
        } catch (Exception e) {
            logger.error("缓存埋点数据失败", e);
        }
    }
    /**
     * 获取埋点数据统计
     */
    @Override
    public EventStats getEventStats(String eventType) {
        // 实现统计逻辑
        return trackEventRepository.countByEventType(eventType);
    }
    @Override
    public List<EventStats> getEventStatsByTimeRange(LocalDateTime start, LocalDateTime end) {
        return trackEventRepository.getStatsBetweenTime(start, end);
    }
}

KafkaProducerService.java

package com.analytics.service;
import com.analytics.model.TrackEvent;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;
@Service
public class KafkaProducerService {
    private static final Logger logger = LoggerFactory.getLogger(KafkaProducerService.class);
    private static final String TOPIC = "track-events";
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;
    @Autowired
    private ObjectMapper objectMapper;
    public void sendTrackEvent(TrackEvent event) {
        try {
            String message = objectMapper.writeValueAsString(event);
            kafkaTemplate.send(TOPIC, event.getEventType(), message)
                    .addCallback(
                        result -> logger.debug("Kafka发送成功: {}", event.getEventId()),
                        ex -> logger.error("Kafka发送失败: {}", ex.getMessage())
                    );
        } catch (Exception e) {
            logger.error("序列化埋点数据失败", e);
        }
    }
}

Repository

TrackEventRepository.java

package com.analytics.repository;
import com.analytics.model.TrackEvent;
import com.analytics.model.EventStats;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param;
import java.time.LocalDateTime;
import java.util.List;
public interface TrackEventRepository extends JpaRepository<TrackEvent, Long> {
    long countByEventType(String eventType);
    List<TrackEvent> findByEventTypeAndTimestampBetween(String eventType, Long start, Long end);
    void deleteByCreatedAtBefore(LocalDateTime date);
    @Query("SELECT new com.analytics.model.EventStats(t.eventType, COUNT(t), t.createdAt) " +
           "FROM TrackEvent t " +
           "WHERE t.createdAt BETWEEN :start AND :end " +
           "GROUP BY t.eventType, t.createdAt")
    List<EventStats> getStatsBetweenTime(@Param("start") LocalDateTime start, 
                                        @Param("end") LocalDateTime end);
    @Query("SELECT t.eventName, COUNT(t) as count " +
           "FROM TrackEvent t " +
           "WHERE t.eventType = :eventType AND t.createdAt BETWEEN :start AND :end " +
           "GROUP BY t.eventName")
    List<Object[]> countByEventNameAndTimeRange(...);
}

工具类

IpUtils.java

package com.analytics.utils;
import javax.servlet.http.HttpServletRequest;
public class IpUtils {
    public static String getIpAddr(HttpServletRequest request) {
        String ip = request.getHeader("X-Forwarded-For");
        if (ip == null || ip.length() == 0 || "unknown".equalsIgnoreCase(ip)) {
            ip = request.getHeader("X-Real-IP");
        }
        if (ip == null || ip.length() == 0 || "unknown".equalsIgnoreCase(ip)) {
            ip = request.getRemoteAddr();
        }
        // 多个IP时取第一个
        if (ip != null && ip.contains(",")) {
            ip = ip.split(",")[0].trim();
        }
        return ip;
    }
}

UserAgentUtils.java

package com.analytics.utils;
import com.analytics.model.TrackEvent;
import eu.bitwalker.useragentutils.UserAgent;
import eu.bitwalker.useragentutils.Browser;
import eu.bitwalker.useragentutils.OperatingSystem;
public class UserAgentUtils {
    public static void parseUserAgent(String userAgentString, TrackEvent event) {
        UserAgent userAgent = UserAgent.parseUserAgentString(userAgentString);
        // 解析浏览器
        Browser browser = userAgent.getBrowser();
        event.setBrowser(browser.getName());
        event.setBrowserVersion(userAgent.getBrowserVersion().getVersion());
        // 解析操作系统
        OperatingSystem os = userAgent.getOperatingSystem();
        event.setOs(os.getName());
        // 解析设备类型
        if (os.isMobileDevice()) {
            event.setDeviceType("mobile");
        } else if (os.isTabletDevice()) {
            event.setDeviceType("tablet");
        } else {
            event.setDeviceType("desktop");
        }
    }
}

配置文件

application.yml

server:
  port: 8080
spring:
  datasource:
    url: jdbc:mysql://localhost:3306/analytics?useUnicode=true&characterEncoding=utf-8
    username: root
    password: 123456
    driver-class-name: com.mysql.cj.jdbc.Driver
  jpa:
    hibernate:
      ddl-auto: update
    show-sql: true
    properties:
      hibernate:
        dialect: org.hibernate.dialect.MySQL8Dialect
  redis:
    host: localhost
    port: 6379
    password: 
    timeout: 3000ms
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
    consumer:
      group-id: track-events-group
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
# 异步线程池配置
async:
  executor:
    core-pool-size: 5
    max-pool-size: 20
    queue-capacity: 1000
    thread-name-prefix: track-async-

前端埋点代码示例

<!-- 埋点SDK示例 -->
<script>
// 分析工具SDK
(function(window, document) {
    class AnalyticsTracker {
        constructor() {
            this.serverUrl = 'http://localhost:8080/api/track';
            this.sessionId = this.getSessionId();
        }
        // 跟踪事件
        track(eventType, eventName, data = {}) {
            const payload = {
                eventType: eventType,
                eventName: eventName,
                eventData: data,
                pageUrl: window.location.href,
                pageTitle: document.title,
                referrer: document.referrer,
                timestamp: Date.now(),
                sessionId: this.sessionId,
                userId: this.getUserId()
            };
            // 发送埋点数据
            navigator.sendBeacon(this.serverUrl + '/collect', 
                new Blob([JSON.stringify(payload)], {type: 'application/json'}));
        }
        // 点击事件
        trackClick(element, eventName = 'click') {
            const data = {
                elementId: element.id,
                elementClass: element.className,
                elementText: element.innerText || element.textContent,
                elementTag: element.tagName
            };
            this.track('click', eventName, data);
        }
        // 页面浏览
        trackPageView() {
            this.track('view', 'page_view', {
                pagePath: window.location.pathname
            });
        }
        // 自定义事件
        trackCustomEvent(eventName, data) {
            this.track('custom', eventName, data);
        }
        getSessionId() {
            let sessionId = sessionStorage.getItem('session_id');
            if (!sessionId) {
                sessionId = 'sess_' + Date.now() + '_' + 
                           Math.random().toString(36).substr(2, 9);
                sessionStorage.setItem('session_id', sessionId);
            }
            return sessionId;
        }
        getUserId() {
            // 从cookie或localStorage获取用户ID
            return localStorage.getItem('user_id') || 'anonymous';
        }
    }
    // 全局实例
    window.Analytics = new AnalyticsTracker();
    // 自动监听页面加载
    document.addEventListener('DOMContentLoaded', function() {
        window.Analytics.trackPageView();
    });
    // 监听点击事件
    document.addEventListener('click', function(e) {
        const element = e.target;
        if (element && element.getAttribute('data-track')) {
            window.Analytics.trackClick(element);
        }
    });
})(window, document);
</script>

数据库表结构

-- 创建数据库
CREATE DATABASE IF NOT EXISTS analytics 
DEFAULT CHARACTER SET utf8mb4 
COLLATE utf8mb4_unicode_ci;
USE analytics;
-- 埋点事件表
CREATE TABLE IF NOT EXISTS track_events (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    event_id VARCHAR(64) NOT NULL UNIQUE,
    event_type VARCHAR(32) NOT NULL,
    event_name VARCHAR(128),
    user_id VARCHAR(64),
    device_id VARCHAR(64),
    page_url TEXT,
    page_title VARCHAR(256),
    element_id VARCHAR(128),
    element_class VARCHAR(256),
    element_text TEXT,
    business_data TEXT,
    ip VARCHAR(64),
    user_agent TEXT,
    device_type VARCHAR(16),
    os VARCHAR(64),
    browser VARCHAR(64),
    browser_version VARCHAR(32),
    referrer TEXT,
    session_id VARCHAR(64),
    timestamp BIGINT,
    server_time DATETIME,
    client_time DATETIME,
    created_at DATETIME,
    INDEX idx_event_type (event_type),
    INDEX idx_event_name (event_name),
    INDEX idx_timestamp (timestamp),
    INDEX idx_created_at (created_at),
    INDEX idx_user_id (user_id),
    INDEX idx_session_id (session_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 创建统计索引
CREATE INDEX idx_type_time ON track_events(event_type, timestamp);
CREATE INDEX idx_name_time ON track_events(event_name, timestamp);

数据统计接口示例

package com.analytics.controller;
@RestController
@RequestMapping("/api/analytics")
public class AnalyticsController {
    @Autowired
    private TrackEventRepository repository;
    /**
     * 获取事件统计
     */
    @GetMapping("/stats")
    public Map<String, Object> getStats(@RequestParam String eventType,
                                       @RequestParam(required = false) String startTime,
                                       @RequestParam(required = false) String endTime) {
        Map<String, Object> result = new HashMap<>();
        // 统计总数
        long totalCount = repository.countByEventType(eventType);
        result.put("totalCount", totalCount);
        // 按天统计
        List<Map<String, Object>> dailyStats = repository.getDailyStats(
            eventType, startTime, endTime);
        result.put("dailyStats", dailyStats);
        return result;
    }
    /**
     * 用户行为分析
     */
    @GetMapping("/user")
    public Map<String, Object> getUserAnalysis(@RequestParam String userId) {
        Map<String, Object> result = new HashMap<>();
        // 用户访问频次
        long visitCount = repository.countByUserId(userId);
        result.put("visitCount", visitCount);
        // 获取用户最近行为
        List<TrackEvent> recentActions = repository.getRecentActions(userId);
        result.put("recentActions", recentActions);
        return result;
    }
}

这个完整的埋点数据收集系统案例包含了:

  1. 数据接收:支持GET、POST、批量提交
  2. 数据处理:异步处理、Kafka实时传输
  3. 数据存储:MySQL持久化、Redis缓存
  4. 数据解析:User-Agent解析、IP获取
  5. 前端SDK:JavaScript自动采集
  6. 数据统计:EventStats、用户行为等

您可以根据实际需求调整数据库、消息队列等技术栈,并根据业务场景增加具体的数据分析功能。

抱歉,评论功能暂时关闭!