Dagster 如何用自定义 Compute Log Manager 对计算日志中的 PII 做脱敏?
Dagster 如何用自定义 Compute Log Manager 对计算日志中的 PII 做脱敏【免费下载链接】dagsterAn orchestration platform for the development, production, and observation of data assets.项目地址: https://gitcode.com/GitHub_Trending/da/dagster当 Dagster 资产处理客户数据时调试输出或错误信息可能把邮箱、电话、SSN、信用卡号等 PII 写进 stdout/stderr而被捕获到计算日志compute logs中对任何能访问 Dagster UI 的人可见。官方最佳实践文档 PII redaction in compute logs 给出的解法是自定义一个 Compute Log Manager在日志进入 UI 之前用[PII 类型]标签自动替换敏感内容。文档提供了两种实现——读取时脱敏redact on read和写入时脱敏redact on write区别在于磁盘上的原始日志是否保留。先看文档中用于演示的资产example_asset.pyimport dagster as dg dg.asset def process_customer_data() - dg.MaterializeResult: # This log output contains PII that will be redacted print(Processing data for customer: John Doe) # noqa: T201 print(Email: john.doeexample.com) # noqa: T201 print(Phone: 555-123-4567) # noqa: T201 print(SSN: 123-45-6789) # noqa: T201 print(Credit Card: 4111-1111-1111-1111) # noqa: T201 print(IP Address: 192.168.1.100) # noqa: T201 return dg.MaterializeResult(metadata{status: processed})不脱敏时日志会把敏感信息完整展示出来以下内容为文档示例Processing data for customer: John Doe Email: john.doeexample.com Phone: 555-123-4567 SSN: 123-45-6789 Credit Card: 4111-1111-1111-1111 IP Address: 192.168.1.100准备共享的 PII 脱敏函数两种方案共用同一个基于正则的脱敏函数pii_redactor.py把匹配到的 PII 替换为[类型名]标签import re # Define regex patterns for common PII types PII_PATTERNS { EMAIL: r\b[A-Za-z0-9._%-][A-Za-z0-9.-]\.[A-Z|a-z]{2,}\b, PHONE: r\b\d{3}[-.]?\d{3}[-.]?\d{4}\b, SSN: r\b\d{3}-\d{2}-\d{4}\b, CREDIT_CARD: r\b\d{4}[-\s]?\d{4}[-\s]?\d{4}[-\s]?\d{4}\b, IP_ADDRESS: r\b(?:\d{1,3}\.){3}\d{1,3}\b, } def redact_pii(text: str) - str: Redact PII from text using regex patterns. redacted text for pii_type, pattern in PII_PATTERNS.items(): redacted re.sub(pattern, f[{pii_type}], redacted) return redacted例如匹配到邮箱地址后该行会被替换为[EMAIL]。需要覆盖更多类型时按同样的模式往PII_PATTERNS里加条目即可。方案一读取时脱敏Redact on read这个方案在日志被读取用于展示时做脱敏磁盘上保留未脱敏的原始日志适合管理员需要访问原始日志做调试或审计的场景。代价是文档指出云部署时文件必须在上传前额外处理脱敏。自定义 Manager 继承LocalComputeLogManager并实现ConfigurableClass重写读取日志的两个入口完整实现见 pii_compute_log_manager.pyfrom collections.abc import Sequence from dagster import Bool, Field, StringSource from dagster._core.storage.captured_log_manager import ( # ty: ignore[unresolved-import] CapturedLogData, ) from dagster._core.storage.compute_log_manager import ComputeIOType from dagster._core.storage.local_compute_log_manager import LocalComputeLogManager from dagster._serdes import ConfigurableClass, ConfigurableClassData from .pii_redactor import redact_pii class PIIComputeLogManager(LocalComputeLogManager, ConfigurableClass): A compute log manager that redacts PII from logs before displaying them. def __init__( self, base_dir: str compute_logs, redact_for_ui: bool True, inst_data: ConfigurableClassData | None None, ): super().__init__(base_dir) self.redact_for_ui redact_for_ui self._inst_data inst_data property def inst_data(self) - ConfigurableClassData | None: return self._inst_data classmethod def config_type(cls): return { base_dir: Field( StringSource, default_valuecompute_logs, descriptionBase directory for storing compute logs, ), redact_for_ui: Field( Bool, default_valueTrue, descriptionWhether to redact PII for UI display, ), } classmethod def from_config_value(cls, inst_data: ConfigurableClassData | None, config_value): return cls(inst_datainst_data, **config_value) def _redact_bytes(self, data: bytes | None) - bytes | None: Apply PII redaction to bytes data. if not data or not self.redact_for_ui: return data text data.decode(utf-8, errorsreplace) redacted_text redact_pii(text) return redacted_text.encode(utf-8) def get_log_data( self, log_key: Sequence[str], cursor: str | None None, max_bytes: int | None None, ) - CapturedLogData: Override to apply PII redaction when logs are read. original_data super().get_log_data(log_key, cursor, max_bytes) if not self.redact_for_ui: return original_data return CapturedLogData( log_keyoriginal_data.log_key, stdoutself._redact_bytes(original_data.stdout), stderrself._redact_bytes(original_data.stderr), cursororiginal_data.cursor, ) def get_log_data_for_type( self, log_key: Sequence[str], io_type: ComputeIOType, offset: int | None 0, max_bytes: int | None None, ) - tuple[bytes | None, int]: Override to apply PII redaction to log data chunks. data, new_offset super().get_log_data_for_type(log_key, io_type, offset, max_bytes) if self.redact_for_ui and data: return self._redact_bytes(data), new_offset return data, new_offset关键点是两处重写get_log_data覆盖完整日志的读取get_log_data_for_type覆盖按流stdout/stderr分块读取两者的数据都经过_redact_bytes处理。然后在dagster.yaml的compute_logs段注册该配置段的通用写法见 dagster-yaml.md默认的LocalComputeLogManager被替换为自定义类compute_logs: module: your_project.pii_compute_log_manager class: PIIComputeLogManager config: base_dir: compute_logs redact_for_ui: true其中your_project是文档中的占位写法需替换为你自己项目中包含该模块的 Python 包路径保证module可被导入base_dir与redact_for_ui都是config_type里声明的配置项。验证方式跑一次会打印 PII 的资产如上面的process_customer_data在 UI 的计算日志中应看到[EMAIL]、[SSN]等标签而非原始值同时直接查看base_dir下的日志文件磁盘上仍是未脱敏的原文——这正是该方案保留原始日志的设计。如果 UI 里还显示明文先确认redact_for_ui是否为true以及module/class是否指向了你的自定义类。方案二写入时脱敏Redact on write如果对安全要求更高可以在日志写入磁盘时就完成脱敏敏感数据永远不落盘云部署也无需额外处理。代价是文档明确说明的权衡——原始数据一旦写入即无法恢复不能用于事后调试。实现上重写capture_logs用PIIRedactingStream包住 stdout/stderr 流每段写入的数据先过一遍redact_pii完整实现见 pii_log_manager_write.pyfrom collections.abc import Iterator, Sequence from contextlib import contextmanager from dagster import Bool, Field, StringSource from dagster._core.storage.captured_log_manager import ( # ty: ignore[unresolved-import] CapturedLogContext, ) from dagster._core.storage.local_compute_log_manager import LocalComputeLogManager from dagster._serdes import ConfigurableClass, ConfigurableClassData from .pii_redactor import redact_pii class PIIRedactingStream: A stream wrapper that redacts PII as data is written. def __init__(self, original_stream, encoding: str utf-8): self.original original_stream self.encoding encoding def write(self, data): if isinstance(data, bytes): text data.decode(self.encoding, errorsreplace) redacted redact_pii(text) return self.original.write(redacted.encode(self.encoding)) elif isinstance(data, str) and data.strip(): redacted redact_pii(data) return self.original.write(redacted) return self.original.write(data) def flush(self): return self.original.flush() def fileno(self): return self.original.fileno() class PIIComputeLogManagerWrite(LocalComputeLogManager, ConfigurableClass): A compute log manager that redacts PII when logs are written to disk. def __init__( self, base_dir: str compute_logs, redact_on_write: bool True, inst_data: ConfigurableClassData | None None, ): super().__init__(base_dir) self.redact_on_write redact_on_write self._inst_data inst_data property def inst_data(self) - ConfigurableClassData | None: return self._inst_data classmethod def config_type(cls): return { base_dir: Field( StringSource, default_valuecompute_logs, descriptionBase directory for storing compute logs, ), redact_on_write: Field( Bool, default_valueTrue, descriptionWhether to redact PII when writing logs, ), } classmethod def from_config_value(cls, inst_data: ConfigurableClassData | None, config_value): return cls(inst_datainst_data, **config_value) contextmanager def capture_logs(self, log_key: Sequence[str]) - Iterator[CapturedLogContext]: Override capture_logs to wrap streams with PII redaction. with super().capture_logs(log_key) as context: if self.redact_on_write: # Wrap the output streams with PII-redacting wrappers original_out context.out_fd # ty: ignore[unresolved-attribute] original_err context.err_fd # ty: ignore[unresolved-attribute] context._out_fd PIIRedactingStream(original_out) # noqa: SLF001 # ty: ignore[unresolved-attribute] context._err_fd PIIRedactingStream(original_err) # noqa: SLF001 # ty: ignore[unresolved-attribute] yield context配置同样是改dagster.yamlcompute_logs: module: your_project.pii_log_manager_write class: PIIComputeLogManagerWrite config: base_dir: compute_logs redact_on_write: true验证方式运行同一资产后UI 日志与磁盘日志应一致——都只有[EMAIL]等脱敏标签。直接打开base_dir下对应的日志文件确认其中不存在明文 PII即说明脱敏发生在落盘之前。两种方案怎么选文档给出的对照方案适用场景读取时脱敏磁盘保留原始日志用于调试UI 中脱敏写入时脱敏最高安全性PII 不落盘云部署开箱即用如果合规要求是原始日志可供审计选方案一如果要求是PII 绝对不进存储选方案二并接受原文不可恢复。可选用 Microsoft Presidio 替代正则检测文档同时提到对合规要求更严格的部署可以考虑用 ML 方式的 Microsoft Presidio 做 PII 检测它的准确率更高且支持 30 实体类型含国际格式。做法是把上面基于正则的redact_pii函数替换为from presidio_analyzer import AnalyzerEngine from presidio_anonymizer import AnonymizerEngine analyzer AnalyzerEngine() anonymizer AnonymizerEngine() def redact_pii(text: str) - str: results analyzer.analyze(texttext, languageen) anonymized anonymizer.anonymize(texttext, analyzer_resultsresults) return anonymized.text两个 Compute Log Manager 的实现无需改动因为都只依赖redact_pii(text) - str这一个入口。【免费下载链接】dagsterAn orchestration platform for the development, production, and observation of data assets.项目地址: https://gitcode.com/GitHub_Trending/da/dagster创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考