收集 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 存储桶

  1. 前往 Google Cloud 控制台
  2. 选择您的项目或创建新项目。
  3. 在导航菜单中,依次前往 Cloud Storage > 存储分区
  4. 点击创建存储分区
  5. 提供以下配置详细信息:

    设置
    为存储桶命名 输入一个全局唯一的名称(例如 netflow-logs-bucket
    位置类型 根据您的需求进行选择(区域级、双区域、多区域)
    位置 选择营业地点(例如 us-central1
    存储类别 标准(建议用于经常访问的日志)
    访问权限控制 均匀(推荐)
    保护工具 可选:启用对象版本控制或保留政策
  6. 点击创建

为 Cloud Run 函数创建服务账号

Cloud Run 函数需要一个服务账号,该账号具有写入 GCS 存储桶的权限,并且可以由 Pub/Sub 调用。

创建服务账号

  1. GCP 控制台中,依次前往 IAM 和管理 > 服务账号
  2. 点击创建服务账号
  3. 提供以下配置详细信息:
    • 服务账号名称:输入 netflow-collector-sa
    • 服务账号说明:输入 Service account for Cloud Run function to collect netflow logs
  4. 点击创建并继续
  5. 向此服务账号授予对项目的访问权限部分中,添加以下角色:
    1. 点击选择角色
    2. 搜索并选择 Storage Object Admin
    3. 点击 + 添加其他角色
    4. 搜索并选择 Cloud Run Invoker
    5. 点击 + 添加其他角色
    6. 搜索并选择 Cloud Functions Invoker
  6. 点击继续
  7. 点击完成

必须拥有这些角色,才能:

  • Storage Object Admin:将日志写入 GCS 存储桶并管理状态文件
  • Cloud Run Invoker:允许 Pub/Sub 调用函数
  • Cloud Functions Invoker:允许调用函数

授予对 GCS 存储桶的 IAM 权限

向服务账号授予对 GCS 存储桶的写入权限:

  1. 前往 Cloud Storage > 存储分区
  2. 点击您的存储桶名称(例如 netflow-logs-bucket)。
  3. 前往权限标签页。
  4. 点击授予访问权限
  5. 提供以下配置详细信息:
    • 添加主账号:输入服务账号电子邮件地址(例如 netflow-collector-sa@PROJECT_ID.iam.gserviceaccount.com
    • 分配角色:选择 Storage Object Admin
  6. 点击保存

创建 Pub/Sub 主题

创建一个 Pub/Sub 主题,Cloud Scheduler 将向该主题发布消息,而 Cloud Run 函数将订阅该主题。

  1. GCP 控制台中,前往 Pub/Sub > 主题
  2. 点击创建主题
  3. 提供以下配置详细信息:
    • 主题 ID:输入 netflow-trigger
    • 将其他设置保留为默认值
  4. 点击创建

创建 Cloud Run 函数以收集 NetFlow 日志

Cloud Run 函数将运行一个 OpenTelemetry 收集器,该收集器接收 NetFlow 数据并将其导出到 GCS。

  1. GCP 控制台中,前往 Cloud Run
  2. 点击创建服务
  3. 选择函数(使用内嵌编辑器创建函数)。
  4. 配置部分中,提供以下配置详细信息:

    设置
    Service 名称 netflow-collector
    区域 选择与您的 GCS 存储桶匹配的区域(例如 us-central1
    运行时 选择 Python 3.12 或更高版本
  5. 触发器(可选)部分中:

    1. 点击 + 添加触发器
    2. 选择 Cloud Pub/Sub
    3. 选择 Cloud Pub/Sub 主题部分,选择 Pub/Sub 主题 netflow-trigger
    4. 点击保存
  6. 身份验证部分中:

    1. 选择需要进行身份验证
    2. 检查 Identity and Access Management (IAM)
  7. 向下滚动并展开容器、网络、安全性

  8. 前往安全性标签页:

    • 服务账号:选择服务账号 netflow-collector-sa
  9. 前往容器标签页:

    1. 点击变量和密钥
    2. 为每个环境变量点击+ 添加变量
    变量名称 示例值 说明
    GCS_BUCKET netflow-logs-bucket GCS 存储桶名称
    GCS_PREFIX netflow 日志文件的前缀
    NETFLOW_PORT 2055 NetFlow 接收器的端口
    NETFLOW_SCHEME netflow 方案类型:netflow、sflow 或 ipfix
    NETFLOW_SOCKETS 4 UDP 套接字数量
    NETFLOW_WORKERS 8 解码器工作器数量
  10. 变量和 Secret 部分中,向下滚动到请求

    • 请求超时:输入 600 秒(10 分钟)
  11. 前往设置标签页:

    • 资源部分中:
      • 内存:选择 1 GiB 或更高
      • CPU:选择 2
  12. 修订版本伸缩部分中:

    • 实例数下限:输入 1
    • 实例数上限:输入 10
  13. 点击创建

  14. 等待服务创建完成(1-2 分钟)。

  15. 创建服务后,系统会自动打开内嵌代码编辑器

添加函数代码

  1. 入口点字段中输入 main
  2. 在内嵌代码编辑器中,创建三个文件:

    • 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. 点击部署以保存并部署该函数。

  4. 等待部署完成(3-5 分钟)。

配置网络设备以发送 NetFlow 数据

将网络设备(路由器、交换机、防火墙)配置为将 NetFlow、sFlow 或 IPFIX 数据导出到 Cloud Run 函数。

获取 Cloud Run 函数的外部 IP

  1. GCP Console 中,前往 Cloud Run > 服务
  2. 点击函数名称 netflow-collector
  3. 复制顶部显示的 网址(例如 https://netflow-collector-xxxxx-uc.a.run.app)。
  4. 从网址中提取主机名。
  5. 使用 nslookupdig 解析 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 函数。

  1. GCP Console 中,前往 Cloud Scheduler
  2. 点击创建作业
  3. 提供以下配置详细信息:

    设置
    名称 netflow-collector-hourly
    区域 选择与 Cloud Run 函数相同的区域
    频率 0 * * * *(每小时一次,整点时)
    时区 选择时区(建议选择世界协调时间 [UTC])
    目标类型 Pub/Sub
    主题 选择 Pub/Sub 主题 netflow-trigger
    消息正文 {}(空 JSON 对象)
  4. 点击创建

时间表频率选项

根据日志量和延迟时间要求选择频次:

频率 Cron 表达式 使用场景
每隔 5 分钟 */5 * * * * 大批量、低延迟
每隔 15 分钟 */15 * * * * 搜索量中等
每小时 0 * * * * 标准(推荐)
每 6 小时 0 */6 * * * 低成交量、批处理

测试集成

  1. Cloud Scheduler 控制台中,找到作业 netflow-collector-hourly
  2. 点击强制运行以手动触发作业。
  3. 等待几秒钟。
  4. 前往 Cloud Run > 服务
  5. 点击函数名称 netflow-collector
  6. 点击日志标签页。
  7. 验证函数是否已成功执行。查找:

    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
    
  8. 前往 Cloud Storage > 存储分区

  9. 点击存储桶名称 netflow-logs-bucket

  10. 前往前缀文件夹 netflow/

  11. 验证是否已创建具有当前时间戳的新 .ndjson 文件。

如果您在日志中看到错误,请执行以下操作:

  • 收集器意外退出:检查 NetFlow 接收器配置
  • 没有要上传的数据:验证网络设备是否正在将 NetFlow 数据发送到正确的 IP 和端口
  • 权限遭拒:检查服务账号是否具有 Storage Object Admin 角色
  • 端口绑定错误:确保端口 2055 未被使用

检索 Google SecOps 服务账号

Google SecOps 使用唯一的服务账号从您的 GCS 存储桶中读取数据。您必须授予此服务账号对您的存储桶的访问权限。

获取服务账号电子邮件地址

  1. 依次前往 SIEM 设置 > Feed
  2. 点击添加新 Feed
  3. 点击配置单个 Feed
  4. Feed 名称字段中,输入 Feed 的名称(例如 OpenTelemetry Netflow Logs)。
  5. 选择 Google Cloud Storage V2 作为来源类型
  6. 选择 NETFLOW_OTEL 作为日志类型
  7. 点击获取服务账号。系统会显示一个唯一的服务账号电子邮件地址。例如:

    chronicle-12345678@chronicle-gcp-prod.iam.gserviceaccount.com
    
  8. 复制此电子邮件地址,以便在下一步中使用。

  9. 点击下一步

  10. 为以下输入参数指定值:

    • 存储桶网址:输入带有前缀路径的 GCS 存储桶 URI:

      gs://netflow-logs-bucket/netflow/
      
    • 来源删除选项:根据您的偏好选择删除选项:

      • 永不:转移后永不删除任何文件(建议用于测试)。
      • 删除已转移的文件:在成功转移后删除文件。
      • 删除已转移的文件和空目录:成功转移后删除文件和空目录。

    • 文件存在时间上限:包含在过去指定天数内修改的文件(默认值为 180 天)

    • 资产命名空间资产命名空间

    • 注入标签:要应用于此 Feed 中事件的标签

  11. 点击下一步

  12. 最终确定界面中查看新的 Feed 配置,然后点击提交

向 Google SecOps 服务账号授予 IAM 权限

Google SecOps 服务账号需要您的 GCS 存储桶的 Storage Object Viewer 角色。

  1. 前往 Cloud Storage > 存储分区
  2. 点击存储桶名称 netflow-logs-bucket
  3. 前往权限标签页。
  4. 点击授予访问权限
  5. 提供以下配置详细信息:
    • 添加主账号:粘贴 Google SecOps 服务账号电子邮件地址
    • 分配角色:选择 Storage Object Viewer
  6. 点击保存

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 从变更日志映射

更新日志

查看相应解析器的更改日志

需要更多帮助?获得社区成员和 Google SecOps 专业人士的解答。