- 示例工程
【免费下载链接】python-docs-samples
Code samples used on cloud.google.com
导读
本文以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))四个核心步骤:
- 读取:
io.ReadFromPubSub(subscription=...)从订阅持续拉取日志条目消息; - 解析:
PayloadAsJson(一个DoFn)将 Pub/Sub 消息字节解码为 UTF-8 并json.loads成字典; - 窗口聚合:
WindowInto(FixedWindows(window_size))将消息按固定时间窗口分组,默认窗口 60 秒; - 批处理与写日志:
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_subscription | str | 必填 | 订阅资源名,格式projects/<PROJECT_ID>/subscription/<SUBSCRIPTION_ID> |
--destination_log_name | str | 必填 | 目标日志名,格式projects/<PROJECT_ID>/logs/<LOG_ID> |
--window_size | float | 60.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)方法的关键点:
- 构造 Table 结构:DLP 支持以
Table形式批量处理多个内容项。代码把每个日志条目的textPayload字段提取为Row({"values": [{"string_value": payload}]}),表头为textPayload; - 调用 API:
dlp_client.deidentify_content(request={...}),请求中parent形如projects/<PROJECT_ID>/locations/<region>,同时携带inspect_config、deidentify_config与批量table; - 回写脱敏结果:按索引把响应中
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
相关推荐
Dataflow GPU 实战:基于 python-docs-samples 构建并运行 TensorFlow 最小 GPU 流水线
Dataflow GPU 实战:基于 python docs samples 构建并运行 TensorFlow 最小 GPU 流水线 Apache Beam 与
示例工程python-docs-samples 实战:基于 Cloud Run Job 将 Cloud Storage 导出的日志回灌到 Cloud Logging
python docs samples 实战:基于 Cloud Run Job 将 Cloud Storage 导出的日志回灌到 Cloud Logging 本
示例工程Dataflow 自定义容器实战:基于 python-docs-samples 构建 Dataflow Worker 专用镜像
Dataflow 自定义容器实战:基于 python docs samples 构建 Dataflow Worker 专用镜像 本指南以 python docs
示例工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考