AbhiramiRajeev/Ingestion-Service
GitHub: AbhiramiRajeev/Ingestion-Service
一个基于 Go 和 Gin 的安全事件接入服务,通过 HTTP 收集并验证事件数据后将其发布到 Kafka 以供下游处理。
Stars: 0 | Forks: 0
# Ingestion 服务
## 问题陈述
现代安全系统会实时地从多个端点产生大量事件 —— 登录尝试、可疑活动、暴力破解信号等。在不丢失数据或产生瓶颈的情况下,收集、验证并可靠地将这些事件转发到处理流水线至关重要。
该服务通过提供一个轻量级、经过身份验证的接入层来解决此问题,该层通过 HTTP 接收事件,对其进行验证,并立即将其发布到 Kafka,供下游消费者进行处理。
## 架构
```
┌──────────────┐
│ Client / │
│ Agent │
└──────┬───────┘
│ POST /ingest
│ Authorization: Bearer
▼
┌──────────────┐
│ Ingestion │
│ Service │
│ (Gin) │
└──────┬───────┘
│ Validate API Key
│ Bind + Validate Payload
│ Marshal to JSON
▼
┌──────────────┐
│ Sync Kafka │
│ Producer │
└──────┬───────┘
│ Publish to topic
▼
┌──────────────┐
│ Kafka │
│ Broker │
│ (login_events)│
└──────────────┘
```
## 请求流程
```
API Request
│
▼
JWT Validation
│
▼
Validate Payload
│
▼
Publish to Kafka
│
▼
Return 200 OK
```
## API 端点
### `POST /ingest`
接收安全事件并将其发布到 Kafka。
**请求头:**
| Header | 值 | 必填 |
|---------------|------------------------|----------|
| Authorization | `Bearer ` | 是 |
| Content-Type | `application/json` | 是 |
**请求体:**
```
{
"event_type": "login_attempt",
"user_id": "alice",
"ip_address": "1.2.3.4",
"status": "failure",
"timestamp": "2025-06-16T12:00:00Z"
}
```
**响应 (200 OK):**
```
{
"message": "Event ingested successfully"
}
```
**错误响应:**
| 状态码 | Body | 何时触发 |
|--------|-----------------------------------------|----------------------------|
| 401 | `{ "error": "Unauthorized" }` | API key 缺失或无效 |
| 400 | `{ "error": "Invalid request body" }` | 格式错误的 JSON 或缺少必填字段 |
| 500 | `{ "error": "Failed to send message to Kafka" }` | Kafka 发布失败 |
### `GET /health`
健康检查端点。目前为一个存根。
## 身份验证流程
```
Client sends request
│
▼
Extract "Authorization" header
│
▼
Split header into "Bearer" +
│
▼
Compare token against configured API_KEY
│
├── Match → proceed with request
│
└── No match → return 401 Unauthorized
```
身份验证是一个简单的 API key 检查 —— 客户端必须在 `Authorization` 请求头中发送有效的 `API_KEY`(设置为环境变量)。这里没有 JWT 解码;而是直接将原始 token 与该 key 进行比对。
## Kafka 事件格式
事件以 JSON 格式发布到 `login_events` 主题,其 schema 如下:
| 字段 | 类型 | 描述 |
|--------------|--------|-------------------------------------------|
| `event_type` | string | 事件类型(例如 `login_attempt`) |
| `user_id` | string | 相关用户的标识符 |
| `ip_address` | string | 事件的源 IP 地址 |
| `status` | string | 结果(例如 `success`, `failure`) |
| `timestamp` | string | 事件的 ISO 8601 时间戳 |
**Kafka 消息示例:**
```
{
"event_type": "login_attempt",
"user_id": "alice",
"ip_address": "1.2.3.4",
"status": "failure",
"timestamp": "2025-06-16T12:00:00Z"
}
```
## 目录结构
```
.
├── cmd/
│ └── main.go # Entrypoint — wires Kafka, Gin router, and handlers
├── config/
│ ├── config.go # Config struct definition
│ └── config.yaml # Default configuration values
├── internal/
│ ├── auth.go # API key validation logic
│ ├── handler.go # HTTP handlers (IngestEvent, GetHealth)
│ └── kafka.go # Kafka producer initialization
├── models/
│ └── models.go # Event data model
├── Dockerfile # Multi-stage build for production image
├── docker-compose.yml # Zookeeper + Kafka for local dev
├── go.mod
└── go.sum
```
## 配置
配置通过 `config/config.yaml` 加载:
```
kafka:
brokers:
- "localhost:9092"
topic:
- "login_events"
auth:
api_key:
- "my-api-key"
server:
port:
- "8080"
```
| 设置项 | 描述 | 默认值 |
|-----------------|------------------------------------------|-----------------|
| `kafka.brokers` | Kafka broker 地址列表 | `localhost:9092`|
| `kafka.topic` | 用于发布事件的 Kafka 主题 | `login_events` |
| `auth.api_key` | 用于身份验证的允许 API key | — |
| `server.port` | HTTP 服务器监听的端口 | `8080` |
## Docker
### 在本地构建并运行
```
# 构建 image
docker build -t ingestion-service .
# 运行 container
docker run -p 8080:8080 -e API_KEY=my-secret-key ingestion-service
```
### 使用 docker-compose (Kafka + Zookeeper)
```
docker-compose up -d
```
这将启动:
| 服务 | 端口 |
|-------------|------|
| Zookeeper | 2181 |
| Kafka | 9092 |
### 完整的本地设置
```
# 1. 启动 Kafka
docker-compose up -d
# 2. 设置 API key
export API_KEY=my-secret-key
# 3. 运行 service
go run cmd/main.go
```
## 技术栈
| 组件 | 技术 |
|---------------|-----------------------------|
| HTTP 服务器 | [Gin](https://github.com/gin-gonic/gin) |
| Kafka 客户端 | [Sarama](https://github.com/IBM/sarama) |
| 日志 | [klog](https://github.com/kubernetes/klog) |
| 语言 | Go 1.24 |
## 未来改进
- **JWT 验证** — 将简单的 API key 检查替换为 JWT token 验证和基于声明的授权
- **异步 producer** — 切换到异步 Kafka producer 以在负载下获得更高的吞吐量
- **结构化日志** — 在日志中添加请求 ID 和结构化字段
- **指标** — 暴露 Prometheus 指标(请求计数、延迟、Kafka 发布速率)
- **限流** — 添加基于客户端的限流以防止滥用
- **Schema 验证** — 根据 JSON schema 或 protobuf 定义验证事件 payload
- **死信队列 (DLQ)** — 将失败的消息路由到 DLQ 以便稍后重试
- **Kafka consumer** — 构建一个 consumer 服务来处理来自 `login_events` 的事件
- **优雅关闭** — 处理 SIGTERM 以清理进行中的请求并干净地关闭 Kafka
- **集成测试** — 添加带有测试 Kafka broker 的测试覆盖率
标签:API网关, EVTX分析, Gin, Go, Kafka, Ruby工具, SonarQube插件, 安全事件采集, 安全编排与自动化响应, 日志审计, 请求拦截