☰
基于 Dataflow 的日志写入时 SSN 自动脱敏流水线实战(python-docs-samples logging/redaction)
2026/10/4 11:12:28 网站建设 项目流程
  • 示例工程

【免费下载链接】python-docs-samples

Code samples used on cloud.google.com

项目地址:https://gitcode.com/GitHub_Trending/py/python-docs-samples
点击查看免费下载

导读

本文以logging/redaction/目录中的官方日志脱敏示例为核心,完整讲解如何构建一条基于 Apache Beam / Dataflow 的流式日志管道:从 Pub/Sub 订阅读取日志条目、调用 Cloud DLP API 在日志写入目标 Log Bucket 之前检测并掩码美国社会安全号码(SSN),最终将脱敏后的日志写入指定日志。读完本文,你将掌握该示例提供的锅炉板代码与最终版代码的每一处差异、关键配置参数(infoTypes、字符掩码、固定窗口与批处理)、自定义 Dataflow 容器镜像的优化手法,以及运行该管道所需的权限与 CLI 参数。

为什么需要在日志写入时做脱敏

应用日志中常常不经意携带敏感数据,例如用户身份证号、信用卡号、社会安全号码等。若日志在落地存储后再做清洗,敏感信息可能已经进入审计日志、被索引或复制,治理成本极高。本示例采用的策略是在摄取(ingestion)阶段即完成检测与脱敏:日志经由 Pub/Sub 进入 Dataflow 管道,管道在把条目写入目标日志之前,用 Cloud DLP 的deidentify_content接口识别并掩码敏感字段,从而保证进入 Log Bucket 的日志已不含明文敏感数据。

目录文件与职责

logging/redaction/目录共包含 5 个文件,README 中明确列出了各文件的用途:

文件作用
Dockerfile自定义 Dataflow 作业容器,用于省去初始化时安装依赖的时间
log_redaction.py锅炉板代码:实现从 Pub/Sub 到目标日志的流式管道骨架
log_redaction_final.py最终版:在锅炉板基础上补全了日志脱敏所需全部改动
requirements.txt安装 Dataflow 作业环境缺失的依赖组件
README.md示例说明与 Cloud Shell 交互式教程入口

说明:原 README 中的[Dockerfile]、[boilerplate]、[final]、[requirements]链接指向仓库内对应文件,本文均转换为仓库根目录相对路径如上表所示。

锅炉板代码:理解流水线骨架

log_redaction.py是学习的起点,它实现了一条可运行的日志摄取管道,但刻意留有三处TODO占位符,正是最终版需要补齐的关键位置。

数据流拓扑

管道主体在run()函数中按以下变换链串联(log_redaction.py):

pipeline | "Read log entries from Pub/Sub" >> io.ReadFromPubSub(subscription=pubsub_subscription) | "Convert log entry payload to Json" >> ParDo(PayloadAsJson()) | "Aggregate payloads in fixed time intervals" >> WindowInto(FixedWindows(window_size)) # Optimize Google API consumption and avoid possible throttling # by calling APIs for batched data and not per each element | "Batch aggregated payloads" >> CombineGlobally(BatchPayloads()).without_defaults() # TODO: Placeholder for redaction transformation | "Ingest to output log" >> ParDo(IngestLogs(destination_log_name))

四个核心步骤:

  1. 读取:io.ReadFromPubSub(subscription=...)从订阅持续拉取日志条目消息;
  2. 解析:PayloadAsJson(一个DoFn)将 Pub/Sub 消息字节解码为 UTF-8 并json.loads成字典;
  3. 窗口聚合:WindowInto(FixedWindows(window_size))将消息按固定时间窗口分组,默认窗口 60 秒;
  4. 批处理与写日志:CombineGlobally(BatchPayloads())把窗口内所有负载聚合成一个列表,IngestLogs一次性批量调用 Cloud Logging 写入接口。

批量聚合的设计意图

BatchPayloads继承自 Beam 的CombineFn(log_redaction.py),实现create_accumulator、add_input、merge_accumulators、extract_output四个方法,将窗口内元素累积到列表。代码注释点明了动机:"Optimize Google API consumption and avoid possible throttling by calling APIs for batched data and not per each element"——即在批数据上调用 API,而不是逐条调用,以优化 Google API 消耗并规避限流。最终版的 DLP 调用同样复用该批量设计。

日志写入实现

IngestLogs在setup()中懒初始化logging_v2.Client()与目标 Logger 对象(log_redaction.py),失败时通过logging.error记录并抛出PipelineError。_replace_log_name会把每条日志条目的logName字段改写为目标 Logger 的名字,随后通过self.logger.client.logging_api.write_entries(logs)批量写入。

命令行参数

__main__通过argparse解析三个业务参数(log_redaction.py):

参数类型默认值说明
--pubsub_subscriptionstr必填订阅资源名,格式projects/<PROJECT_ID>/subscription/<SUBSCRIPTION_ID>
--destination_log_namestr必填目标日志名,格式projects/<PROJECT_ID>/logs/<LOG_ID>
--window_sizefloat60.0输出窗口大小(秒)

parse_known_args将剩余参数交给PipelineOptions处理(如--runner=DataflowRunner、--region、--project等 Dataflow 原生选项),并设置streaming=True、save_main_session=True。

锅炉板中另外两处TODO为:# TODO: Place inspection and de-identification configurations(配置占位)与# TODO: Read job's deployment region(区域读取占位),这两处正是最终版要解决的核心问题。

最终版代码:补全脱敏能力

log_redaction_final.py在锅炉板基础上引入google.cloud.dlp_v2,并新增三块内容:脱敏配置、LogRedaction变换、作业区域解析。

检测与脱敏配置

文件顶部定义了两个配置字典(log_redaction_final.py):

INSPECT_CFG = {"info_types": [{"name": "US_SOCIAL_SECURITY_NUMBER"}]} REDACTION_CFG = { "info_type_transformations": { "transformations": [ { "primitive_transformation": { "character_mask_config": {"masking_character": "#"} } } ] } }
  • INSPECT_CFG指定检测的信息类型为US_SOCIAL_SECURITY_NUMBER。DLP 支持大量内置 infoType(如EMAIL_ADDRESS、CREDIT_CARD_NUMBER、PHONE_NUMBER等),只需替换info_types列表即可扩展脱敏范围,配置格式对应 DLPInspectConfigREST 结构;
  • REDACTION_CFG采用character_mask_config字符掩码,把命中的敏感字符替换为#。此处未设置number_to_mask,意味着掩码全部命中字符;若需限制掩码数量,可参考仓库中 DLP 掩码示例 dlp/snippets/Deidentify/deidentify_masking.py 的number_to_mask参数用法。掩码格式对应 DLPDeidentifyTemplate.InfoTypeTransformations结构。

LogRedaction 变换:核心脱敏逻辑

LogRedaction是一个DoFn(log_redaction_final.py),构造时接收region与project_id,在setup()中创建dlp_v2.DlpServiceClient()并做非空校验。

其process(logs)方法的关键点:

  1. 构造 Table 结构:DLP 支持以Table形式批量处理多个内容项。代码把每个日志条目的textPayload字段提取为Row({"values": [{"string_value": payload}]}),表头为textPayload;
  2. 调用 API:dlp_client.deidentify_content(request={...}),请求中parent形如projects/<PROJECT_ID>/locations/<region>,同时携带inspect_config、deidentify_config与批量table;
  3. 回写脱敏结果:按索引把响应中response.item.table.rows[index].values[0].string_value赋回log["textPayload"],即日志主体被替换为脱敏版本。

代码注释还给出一个实战提醒:如果目标项目已存在同名日志副本,可能需要修改insertId(例如log['insertId'] = 'deid-' + log['insertId'])以避免与原始日志冲突。

作业区域解析

锅炉板的# TODO: Read job's deployment region在最终版中落实为(log_redaction_final.py):

region = "us-central1" try: region = pipeline_options.view_as(GoogleCloudOptions).region except AttributeError: pass

默认us-central1,若通过 Dataflow 管道选项传入--region,则优先采用该区域,并将region与从destination_log_name.split("/")[1]解析出的项目 ID 一起传给LogRedaction。

与锅炉板的差异小结

最终版相对锅炉板仅改动/新增四处:引入dlp_v2与GoogleCloudOptions导入、定义INSPECT_CFG/REDACTION_CFG、新增LogRedactionDoFn、在管道变换链的批处理与写日志之间插入"Redact SSN info from logs" >> ParDo(LogRedaction(...))步骤。逐文件 diff 即可清晰看到"骨架 → 完整实现"的演进路径。

自定义 Dataflow 容器:Dockerfile 剖析

Dockerfile 的作用是定制 Dataflow 作业容器以节省初始化时间:

FROM apache/beam_python3.9_sdk@sha256:246c4b813c6de8c240b49ed03c426f413f1768321a3c441413031396a08912f9 # Install google-cloud-logging package that is missing in Beam SDK COPY requirements.txt /tmp RUN pip3 install --upgrade pip && pip3 install -r /tmp/requirements.txt && pip3 check
  • 基础镜像锁定为apache/beam_python3.9_sdk的固定 SHA 摘要,保证可复现;
  • 注释点明:google-cloud-logging包在 Beam SDK 中缺失,因此构建阶段即通过requirements.txt预装;
  • pip3 check用于验证依赖完整性。把依赖安装固化进镜像,作业启动时无需在线安装,从而缩短初始化时间。

配套的 requirements.txt 仅一行:google-cloud-logging>=3.4.0。注意最终版代码还使用google-cloud-dlp,运行最终版时需确保该依赖同样可用。

运行管道与前置条件

按 README 说明,运行该示例需要具备以下权限与资源:

  • 有权启用 Google API;
  • 有权开通 Pub/Sub、Dataflow 与 Cloud Storage 资源;
  • 有权调用 DLP API。

典型启动命令(DataflowRunner 流式作业)可概括为:

python log_redaction_final.py \ --pubsub_subscription projects/<PROJECT_ID>/subscription/<SUBSCRIPTION_ID> \ --destination_log_name projects/<PROJECT_ID>/logs/<LOG_ID> \ --window_size 60 \ --runner DataflowRunner \ --project <PROJECT_ID> \ --region us-central1 \ --staging_location gs://<BUCKET>/staging \ --temp_location gs://<BUCKET>/temp \ --requirements_file requirements.txt \ --sdk_container_image gcr.io/<PROJECT_ID>/<custom-image>:<tag>

其中--sdk_container_image指向按上文 Dockerfile 构建的自定义镜像,其余为 Dataflow 标准管道选项(由parse_known_args透传给PipelineOptions)。

交互式教程

如果你拥有 Google Cloud 账号且能访问 GCP 项目,可直接在 Cloud Console 启动交互式教程(README 中的 "Open in Cloud Shell" 按钮)直观体验示例运行效果;教程会引导完成资源开通、作业提交与日志验证。同样地,运行该教程同样需要上述启用 API、开通 Pub/Sub / Dataflow / Cloud Storage 资源以及调用 DLP API 的权限。

可扩展方向

  • 扩大脱敏范围:在INSPECT_CFG["info_types"]中追加多个 infoType(如EMAIL_ADDRESS、CREDIT_CARD_NUMBER),或为不同 infoType 配置不同变换;
  • 调整掩码行为:在character_mask_config中增加number_to_mask,控制每条命中最多掩码的字符数,可参考 deidentify_masking.py 的参数语义;
  • 批量吞吐优化:窗口大小--window_size直接影响 DLP 调用批次规模,可根据日志量与 DLP 配额权衡调整;
  • 日志去重:若目标日志已存在同类条目,可按最终版注释所示改写insertId前缀以区分脱敏版本。

小结

logging/redaction/示例演示了一条完整的"写入时日志脱敏"链路:Pub/Sub 摄取 → Beam 窗口批处理 → DLP 检测掩码 SSN → Cloud Logging 批量写入,并配套了自定义容器与依赖清单。从锅炉板到最终版的渐进式代码设计,特别适合作为学习 Beam 流式管道与 DLP 集成的实战范本。

  • 示例工程

【免费下载链接】python-docs-samples

Code samples used on cloud.google.com

项目地址:https://gitcode.com/GitHub_Trending/py/python-docs-samples
点击查看免费下载
上一篇:终极指南:QUIC流管理实战 - 双向流与单向流应用场景深度解析
下一篇:Playnite:三步跑通多平台游戏库管理

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询