收集 OpenTelemetry Netflow 接收器日志
本文档介绍了如何使用 Google Cloud Storage V2 将 OpenTelemetry Netflow 接收器日志注入到 Google Security Operations。
OpenTelemetry Netflow 接收器是一个开源组件,用于监听 Netflow、sflow 和 IPFIX UDP 流量,并将其转换为 OpenTelemetry 日志记录。这样一来,便可进行网络流量监控和分析,包括协议识别、流量分析、端口使用情况跟踪以及字节/数据包统计。
准备工作
请确保您满足以下前提条件:
- Google SecOps 实例
- 已启用 Cloud Storage API 的 GCP 项目
- 创建和管理 GCS 存储分区的权限
- 管理 GCS 存储分区的 IAM 政策的权限
- 创建 Cloud Run 服务、Pub/Sub 主题和 Cloud Scheduler 作业的权限
- 能够发送 NetFlow、sFlow 或 IPFIX 数据的网络设备
- 拥有配置网络设备导出设置的权限
创建 Google Cloud Storage 存储桶
- 前往 Google Cloud 控制台。
- 选择您的项目或创建新项目。
- 在导航菜单中,依次前往 Cloud Storage > 存储分区。
- 点击创建存储分区。
提供以下配置详细信息:
设置 值 为存储桶命名 输入一个全局唯一的名称(例如 netflow-logs-bucket)位置类型 根据您的需求进行选择(区域级、双区域、多区域) 位置 选择营业地点(例如 us-central1)存储类别 标准(建议用于经常访问的日志) 访问权限控制 均匀(推荐) 保护工具 可选:启用对象版本控制或保留政策 点击创建。
为 Cloud Run 函数创建服务账号
Cloud Run 函数需要一个服务账号,该账号具有写入 GCS 存储桶的权限,并且可以由 Pub/Sub 调用。
创建服务账号
- 在 GCP 控制台中,依次前往 IAM 和管理 > 服务账号。
- 点击创建服务账号。
- 提供以下配置详细信息:
- 服务账号名称:输入
netflow-collector-sa - 服务账号说明:输入
Service account for Cloud Run function to collect netflow logs
- 服务账号名称:输入
- 点击创建并继续。
- 在向此服务账号授予对项目的访问权限部分中,添加以下角色:
- 点击选择角色。
- 搜索并选择 Storage Object Admin。
- 点击 + 添加其他角色。
- 搜索并选择 Cloud Run Invoker。
- 点击 + 添加其他角色。
- 搜索并选择 Cloud Functions Invoker。
- 点击继续。
- 点击完成。
必须拥有这些角色,才能:
- Storage Object Admin:将日志写入 GCS 存储桶并管理状态文件
- Cloud Run Invoker:允许 Pub/Sub 调用函数
- Cloud Functions Invoker:允许调用函数
授予对 GCS 存储桶的 IAM 权限
向服务账号授予对 GCS 存储桶的写入权限:
- 前往 Cloud Storage > 存储分区。
- 点击您的存储桶名称(例如
netflow-logs-bucket)。 - 前往权限标签页。
- 点击授予访问权限。
- 提供以下配置详细信息:
- 添加主账号:输入服务账号电子邮件地址(例如
netflow-collector-sa@PROJECT_ID.iam.gserviceaccount.com) - 分配角色:选择 Storage Object Admin
- 添加主账号:输入服务账号电子邮件地址(例如
- 点击保存。
创建 Pub/Sub 主题
创建一个 Pub/Sub 主题,Cloud Scheduler 将向该主题发布消息,而 Cloud Run 函数将订阅该主题。
- 在 GCP 控制台中,前往 Pub/Sub > 主题。
- 点击创建主题。
- 提供以下配置详细信息:
- 主题 ID:输入
netflow-trigger - 将其他设置保留为默认值
- 主题 ID:输入
- 点击创建。
创建 Cloud Run 函数以收集 NetFlow 日志
Cloud Run 函数将运行一个 OpenTelemetry 收集器,该收集器接收 NetFlow 数据并将其导出到 GCS。
- 在 GCP 控制台中,前往 Cloud Run。
- 点击创建服务。
- 选择函数(使用内嵌编辑器创建函数)。
在配置部分中,提供以下配置详细信息:
设置 值 Service 名称 netflow-collector区域 选择与您的 GCS 存储桶匹配的区域(例如 us-central1)运行时 选择 Python 3.12 或更高版本 在触发器(可选)部分中:
- 点击 + 添加触发器。
- 选择 Cloud Pub/Sub。
- 在选择 Cloud Pub/Sub 主题部分,选择 Pub/Sub 主题
netflow-trigger。 - 点击保存。
在身份验证部分中:
- 选择需要进行身份验证。
- 检查 Identity and Access Management (IAM)。
向下滚动并展开容器、网络、安全性。
前往安全性标签页:
- 服务账号:选择服务账号
netflow-collector-sa
- 服务账号:选择服务账号
前往容器标签页:
- 点击变量和密钥。
- 为每个环境变量点击+ 添加变量:
变量名称 示例值 说明 GCS_BUCKETnetflow-logs-bucketGCS 存储桶名称 GCS_PREFIXnetflow日志文件的前缀 NETFLOW_PORT2055NetFlow 接收器的端口 NETFLOW_SCHEMEnetflow方案类型:netflow、sflow 或 ipfix NETFLOW_SOCKETS4UDP 套接字数量 NETFLOW_WORKERS8解码器工作器数量 在变量和 Secret 部分中,向下滚动到请求:
- 请求超时:输入
600秒(10 分钟)
- 请求超时:输入
前往设置标签页:
- 在资源部分中:
- 内存:选择 1 GiB 或更高
- CPU:选择 2
- 在资源部分中:
在修订版本伸缩部分中:
- 实例数下限:输入
1 - 实例数上限:输入
10
- 实例数下限:输入
点击创建。
等待服务创建完成(1-2 分钟)。
创建服务后,系统会自动打开内嵌代码编辑器。
添加函数代码
- 在入口点字段中输入 main。
在内嵌代码编辑器中,创建三个文件:
- main.py:
import functions_framework from google.cloud import storage import json import os import subprocess import tempfile import signal import time from datetime import datetime, timezone # Initialize Storage client storage_client = storage.Client() # Environment variables GCS_BUCKET = os.environ.get('GCS_BUCKET') GCS_PREFIX = os.environ.get('GCS_PREFIX', 'netflow') NETFLOW_PORT = os.environ.get('NETFLOW_PORT', '2055') NETFLOW_SCHEME = os.environ.get('NETFLOW_SCHEME', 'netflow') NETFLOW_SOCKETS = os.environ.get('NETFLOW_SOCKETS', '4') NETFLOW_WORKERS = os.environ.get('NETFLOW_WORKERS', '8') # Global process handle collector_process = None def create_collector_config(): """Create OpenTelemetry Collector configuration.""" timestamp = datetime.now(timezone.utc).strftime('%Y%m%d_%H%M%S') file_path = f"/tmp/netflow_{timestamp}.ndjson" config = { 'receivers': { 'netflow': { 'scheme': NETFLOW_SCHEME, 'hostname': '0.0.0.0', 'port': int(NETFLOW_PORT), 'sockets': int(NETFLOW_SOCKETS), 'workers': int(NETFLOW_WORKERS), 'queue_size': 5000 } }, 'processors': { 'batch': { 'timeout': '10s', 'send_batch_size': 1000 } }, 'exporters': { 'file': { 'path': file_path, 'format': 'json' } }, 'service': { 'pipelines': { 'logs': { 'receivers': ['netflow'], 'processors': ['batch'], 'exporters': ['file'] } }, 'telemetry': { 'logs': { 'level': 'info' } } } } config_path = '/tmp/otel-config.yaml' with open(config_path, 'w') as f: import yaml yaml.dump(config, f) return config_path, file_path def upload_to_gcs(file_path): """Upload netflow logs to GCS.""" if not os.path.exists(file_path) or os.path.getsize(file_path) == 0: print(f"No data to upload from {file_path}") return bucket = storage_client.bucket(GCS_BUCKET) timestamp = datetime.now(timezone.utc).strftime('%Y%m%d_%H%M%S') object_key = f"{GCS_PREFIX}/netflow_{timestamp}.ndjson" blob = bucket.blob(object_key) blob.upload_from_filename(file_path, content_type='application/x-ndjson') print(f"Uploaded {os.path.getsize(file_path)} bytes to gs://{GCS_BUCKET}/{object_key}") def signal_handler(signum, frame): """Handle shutdown signals.""" global collector_process if collector_process: print("Shutting down collector...") collector_process.terminate() collector_process.wait(timeout=10) @functions_framework.cloud_event def main(cloud_event): """ Cloud Run function to run OpenTelemetry Collector for netflow collection. Args: cloud_event: CloudEvent object containing Pub/Sub message """ global collector_process if not GCS_BUCKET: print('Error: GCS_BUCKET environment variable not set') return try: # Set up signal handlers signal.signal(signal.SIGTERM, signal_handler) signal.signal(signal.SIGINT, signal_handler) # Create collector configuration config_path, file_path = create_collector_config() print(f"Created collector config at {config_path}") print(f"Netflow receiver listening on {NETFLOW_SCHEME}://0.0.0.0:{NETFLOW_PORT}") # Start OpenTelemetry Collector collector_process = subprocess.Popen( ['/otelcol-contrib', '--config', config_path], stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True ) print(f"Started OpenTelemetry Collector (PID: {collector_process.pid})") # Run for collection period (e.g., 5 minutes) collection_time = 300 print(f"Collecting netflow data for {collection_time} seconds...") start_time = time.time() while time.time() - start_time < collection_time: if collector_process.poll() is not None: stdout, stderr = collector_process.communicate() print(f"Collector exited unexpectedly") print(f"STDOUT: {stdout}") print(f"STDERR: {stderr}") break time.sleep(10) # Stop collector if collector_process.poll() is None: print("Stopping collector...") collector_process.terminate() collector_process.wait(timeout=10) # Upload collected data upload_to_gcs(file_path) # Cleanup if os.path.exists(file_path): os.remove(file_path) if os.path.exists(config_path): os.remove(config_path) print("Netflow collection completed successfully") except Exception as e: print(f'Error during netflow collection: {str(e)}') if collector_process and collector_process.poll() is None: collector_process.terminate() raise- requirements.txt:
functions-framework==3.* google-cloud-storage==2.* PyYAML==6.*- Dockerfile:
FROM python:3.12-slim # Install OpenTelemetry Collector Contrib RUN apt-get update && apt-get install -y wget && \ wget https://github.com/open-telemetry/opentelemetry-collector-releases/releases/download/v0.144.0/otelcol-contrib_0.144.0_linux_amd64.deb && \ dpkg -i otelcol-contrib_0.144.0_linux_amd64.deb && \ rm otelcol-contrib_0.144.0_linux_amd64.deb && \ apt-get clean # Set working directory WORKDIR /app # Copy requirements and install COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt # Copy function code COPY main.py . # Expose netflow port EXPOSE 2055/udp # Run function CMD ["functions-framework", "--target=main", "--port=8080"]点击部署以保存并部署该函数。
等待部署完成(3-5 分钟)。
配置网络设备以发送 NetFlow 数据
将网络设备(路由器、交换机、防火墙)配置为将 NetFlow、sFlow 或 IPFIX 数据导出到 Cloud Run 函数。
获取 Cloud Run 函数的外部 IP
- 在 GCP Console 中,前往 Cloud Run > 服务。
- 点击函数名称
netflow-collector。 - 复制顶部显示的 网址(例如
https://netflow-collector-xxxxx-uc.a.run.app)。 - 从网址中提取主机名。
使用
nslookup或dig解析 IP 地址:nslookup netflow-collector-xxxxx-uc.a.run.app
示例:Cisco 路由器 NetFlow 配置
! Configure netflow exporter flow exporter OTEL-EXPORTER destination <CLOUD_RUN_IP> transport udp 2055 source <INTERFACE> export-protocol netflow-v9 ! Configure flow monitor flow monitor OTEL-MONITOR exporter OTEL-EXPORTER record netflow ipv4 original-input ! Apply to interface interface GigabitEthernet0/0 ip flow monitor OTEL-MONITOR input ip flow monitor OTEL-MONITOR output示例:通用 Netflow 配置
对于大多数网络设备,请配置以下设置:
- NetFlow 版本:NetFlow v5、v9 或 IPFIX
- 收集器 IP:Cloud Run 函数 IP 地址
- 收集器端口:
2055(或配置的端口) - 协议:UDP
- 有效超时时间:60 秒(建议值)
- 不活动超时时间:15 秒(建议值)
创建 Cloud Scheduler 作业
Cloud Scheduler 会定期向 Pub/Sub 主题发布消息,从而触发 Cloud Run 函数。
- 在 GCP Console 中,前往 Cloud Scheduler。
- 点击创建作业。
提供以下配置详细信息:
设置 值 名称 netflow-collector-hourly区域 选择与 Cloud Run 函数相同的区域 频率 0 * * * *(每小时一次,整点时)时区 选择时区(建议选择世界协调时间 [UTC]) 目标类型 Pub/Sub 主题 选择 Pub/Sub 主题 netflow-trigger消息正文 {}(空 JSON 对象)点击创建。
时间表频率选项
根据日志量和延迟时间要求选择频次:
| 频率 | Cron 表达式 | 使用场景 |
|---|---|---|
| 每隔 5 分钟 | */5 * * * * |
大批量、低延迟 |
| 每隔 15 分钟 | */15 * * * * |
搜索量中等 |
| 每小时 | 0 * * * * |
标准(推荐) |
| 每 6 小时 | 0 */6 * * * |
低成交量、批处理 |
测试集成
- 在 Cloud Scheduler 控制台中,找到作业
netflow-collector-hourly。 - 点击强制运行以手动触发作业。
- 等待几秒钟。
- 前往 Cloud Run > 服务。
- 点击函数名称
netflow-collector。 - 点击日志标签页。
验证函数是否已成功执行。查找:
Started OpenTelemetry Collector (PID: ...) Netflow receiver listening on netflow://0.0.0.0:2055 Collecting netflow data for 300 seconds... Uploaded ... bytes to gs://netflow-logs-bucket/netflow/netflow_YYYYMMDD_HHMMSS.ndjson Netflow collection completed successfully前往 Cloud Storage > 存储分区。
点击存储桶名称
netflow-logs-bucket。前往前缀文件夹
netflow/。验证是否已创建具有当前时间戳的新
.ndjson文件。
如果您在日志中看到错误,请执行以下操作:
- 收集器意外退出:检查 NetFlow 接收器配置
- 没有要上传的数据:验证网络设备是否正在将 NetFlow 数据发送到正确的 IP 和端口
- 权限遭拒:检查服务账号是否具有 Storage Object Admin 角色
- 端口绑定错误:确保端口 2055 未被使用
检索 Google SecOps 服务账号
Google SecOps 使用唯一的服务账号从您的 GCS 存储桶中读取数据。您必须授予此服务账号对您的存储桶的访问权限。
获取服务账号电子邮件地址
- 依次前往 SIEM 设置 > Feed。
- 点击添加新 Feed。
- 点击配置单个 Feed。
- 在 Feed 名称字段中,输入 Feed 的名称(例如
OpenTelemetry Netflow Logs)。 - 选择 Google Cloud Storage V2 作为来源类型。
- 选择 NETFLOW_OTEL 作为日志类型。
点击获取服务账号。系统会显示一个唯一的服务账号电子邮件地址。例如:
chronicle-12345678@chronicle-gcp-prod.iam.gserviceaccount.com复制此电子邮件地址,以便在下一步中使用。
点击下一步。
为以下输入参数指定值:
存储桶网址:输入带有前缀路径的 GCS 存储桶 URI:
gs://netflow-logs-bucket/netflow/
来源删除选项:根据您的偏好选择删除选项:
- 永不:转移后永不删除任何文件(建议用于测试)。
- 删除已转移的文件:在成功转移后删除文件。
删除已转移的文件和空目录:成功转移后删除文件和空目录。
文件存在时间上限:包含在过去指定天数内修改的文件(默认值为 180 天)
资产命名空间:资产命名空间
注入标签:要应用于此 Feed 中事件的标签
点击下一步。
在最终确定界面中查看新的 Feed 配置,然后点击提交。
向 Google SecOps 服务账号授予 IAM 权限
Google SecOps 服务账号需要您的 GCS 存储桶的 Storage Object Viewer 角色。
- 前往 Cloud Storage > 存储分区。
- 点击存储桶名称
netflow-logs-bucket。 - 前往权限标签页。
- 点击授予访问权限。
- 提供以下配置详细信息:
- 添加主账号:粘贴 Google SecOps 服务账号电子邮件地址
- 分配角色:选择 Storage Object Viewer
- 点击保存。
UDM 映射表
下表显示了 OpenTelemetry Netflow Receiver 日志字段与 Google SecOps UDM 字段的映射关系:
| OpenTelemetry 字段 | UDM 字段 | 说明 |
|---|---|---|
source.address |
principal.ip |
来源 IP 地址 |
source.port |
principal.port |
来源端口号 |
destination.address |
target.ip |
目标 IP 地址 |
destination.port |
target.port |
目标端口号 |
network.transport |
network.ip_protocol |
传输协议(tcp、udp) |
network.type |
network.ip_version |
IP 版本(ipv4、ipv6) |
flow.io.bytes |
network.sent_bytes |
转移的总字节数 |
flow.io.packets |
network.sent_packets |
传输的数据包总数 |
flow.type |
metadata.product_log_id |
流量类型(NETFLOW_V5、NETFLOW_V9、SFLOW_5、IPFIX) |
flow.start |
metadata.event_timestamp |
Flow 开始时间 |
flow.end |
network.session_duration |
Flow 结束时间 |
flow.sampler_address |
intermediary.ip |
Netflow 输出器 IP 地址 |
flow.tcp_flags |
network.tcp_flags |
TCP 标志 |
flow.io.bytes |
event.idm.read_only_udm.network.sent_bytes |
从变更日志映射 |
source.port |
event.idm.read_only_udm.principal.port |
从变更日志映射 |
destination.port |
event.idm.read_only_udm.target.port |
从变更日志映射 |
network.transport |
event.idm.read_only_udm.network.ip_protocol |
从变更日志映射 |
network.type |
event.idm.read_only_udm.security_result.detection_fields |
从变更日志映射 |
flow.sequence_num |
event.idm.read_only_udm.network.session_id |
从变更日志映射 |
flow.sampler_address |
event.idm.read_only_udm.observer.ip |
从变更日志映射 |
flow.type |
event.idm.read_only_udm.metadata.product_event_type |
从变更日志映射 |