本文目录导读:

我将为您提供一个完整的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;
}
}
这个完整的埋点数据收集系统案例包含了:
- 数据接收:支持GET、POST、批量提交
- 数据处理:异步处理、Kafka实时传输
- 数据存储:MySQL持久化、Redis缓存
- 数据解析:User-Agent解析、IP获取
- 前端SDK:JavaScript自动采集
- 数据统计:EventStats、用户行为等
您可以根据实际需求调整数据库、消息队列等技术栈,并根据业务场景增加具体的数据分析功能。