全部文档
文档中心数据流3.0节点组件服务PY代码

PY代码

用于执行输入的 PY 代码,输出结果为非数据集。

输入 Python 代码即可。

提供联想功能并支持变量引用,具体语法参见预置变量,使用 Tab 键接受联想:
1737534034763-01687d04-9c84-4139-bfa1-b787a056a058

若多个PY代码节点都用到相同的方法,推荐在全局设置-PY设置中配置公共脚本,具体配置参见全局设置。

PY 代码适合处理参数、状态和小型结果,不建议承载大批量明细数据。大数据量应尽量保留为数据集,并使用查询、数据转换或 SQL 转换节点处理,避免转换为 Python 列表、字典或 Pandas DataFrame 后造成额外内存占用。更多建议参见并发、队列与内存优化。

案例:多个 PY 节点都使用DataTableMySQL方法。

在全局设置-PY设置中配置公共脚本:
1737534334782-c9600559-ec85-485c-92b1-66bf94dfe416

PY代码节点可以直接调用,无需再重复定义:
1737534384940-1e0d38dd-f64b-4b95-9143-07b3fba796aa

数据流运行引擎基于 Python,PY代码节点除了处理业务参数,也可以调用 Python 标准库查看运行环境、辅助排查问题。可执行的操作取决于节点实际运行环境中安装的库、可访问的目录和进程权限。

在数据流中添加一个PY代码节点,在执行代码中输入:

Copy
import os
dict(os.environ)

执行节点后,返回结果就是当前节点 Python 进程可见的全部环境变量,以字典形式展示。这里查看的是节点实际运行所在容器或进程的环境,不是浏览器所在电脑的环境,也不代表其他 Pod 或整个集群的环境变量。

如果只需确认某个变量,可以按名称查询,例如:

Copy
import os
os.environ.get("TZ")

将 TZ 替换为需要查询的变量名;未配置时返回 None。

环境变量中可能包含数据库密码、访问令牌等敏感信息。排查时优先查询需要的变量,分享节点结果或截图前应脱敏,避免把完整结果写入公开文档或对外响应。

当需要排查 Linux 运行环境中 /dev/shm 下的 sem.mp-* 文件时,可以使用下面的示例扫描文件,并预览满足条件的清理候选项。示例保留了手动删除能力,默认只预览。

  • 在需要排查的 pipeline 运行容器中执行,要求能够访问 Linux /dev/shm 和相关进程的 /proc/<pid>/maps;通过结果中的 hostname 确认实际执行位置。多个 Pod 需要分别确认,单次执行不会检查全部 Pod。

  • 只检查 /dev/shm 下名称以 sem.mp- 开头的条目,实际匹配的是 sem.mp-*,其中 sem 后是英文句点。

  • 已在可读取的进程映射中找到的条目保留;未找到映射,但修改时间未超过阈值的条目也保留。

  • 同时满足“本次未找到映射”和“修改时间早于阈值”的条目进入 would_delete。阈值按当前时间向前推算,1 天为 24 小时;修改时间不等同于最后使用时间。

删除前确认: not_found_in_proc_maps 仅表示本次扫描未发现映射,不等于文件一定没有进程使用。该示例会跳过无法读取的进程信息,权限限制、进程可见范围以及扫描期间的状态变化都可能影响判断。请先使用默认预览模式,由运维核对执行容器、共享内存范围和相关进程状态;暂停相关业务任务并确认候选文件可清理后,再显式开启删除。不要把这段示例直接配置为定时自动清理任务。

添加一个PY代码节点,将以下完整代码粘贴到执行代码中。代码中的 Pipeline.params 读取本次数据流的启动参数,参数配置参见预置变量。

Copy
import os
from collections import defaultdict
from datetime import datetime
from socket import gethostname

SHM_PATH = "/dev/shm"
SEM_PREFIX = "sem.mp-"
MAPS_MARKER = "/dev/shm/sem.mp-"


def _holders_of(names):
    holders = defaultdict(set)
    try:
        pids = os.listdir("/proc")
    except OSError:
        return holders

    for pid in pids:
        if not pid.isdigit():
            continue
        try:
            with open(f"/proc/{pid}/maps", encoding="utf-8", errors="ignore") as f:
                for line in f:
                    if MAPS_MARKER not in line:
                        continue
                    path = line[line.index(MAPS_MARKER):].strip()
                    name = os.path.basename(path).replace(" (deleted)", "").strip()
                    if name in names:
                        holders[name].add(int(pid))
        except (FileNotFoundError, PermissionError, ProcessLookupError, OSError):
            continue
    return holders


param = Pipeline.params or {}
dry_run = bool(param.get("dry_run", True))
min_age_days = float(param.get("min_age_days", 1))
min_age_seconds = min_age_days * 24 * 3600
now = datetime.now().timestamp()
cutoff = now - min_age_seconds

if os.name != "posix" or not os.path.isdir(SHM_PATH):
    return {
        "ok": False,
        "error": "需要 Linux /dev/shm,请在 pipeline 所在 Pod 上跑",
    }

items = []
for entry in os.scandir(SHM_PATH):
    if not entry.name.startswith(SEM_PREFIX):
        continue
    try:
        stat = entry.stat(follow_symlinks=False)
    except FileNotFoundError:
        continue
    items.append({
        "name": entry.name,
        "mtime": stat.st_mtime,
        "time": datetime.fromtimestamp(stat.st_mtime).strftime("%Y-%m-%d %H:%M:%S"),
        "old_enough": stat.st_mtime < cutoff,
    })

names = {item["name"] for item in items}
holders = _holders_of(names)

orphans = []
kept_mapped = []
kept_young = []
for item in items:
    pids = sorted(holders.get(item["name"], []))
    if pids:
        kept_mapped.append(item["name"])
        continue
    if not item["old_enough"]:
        kept_young.append(item["name"])
        continue
    orphans.append(item["name"])

deleted = []
failed = []
if not dry_run:
    for name in orphans:
        path = os.path.join(SHM_PATH, name)
        try:
            os.unlink(path)
            deleted.append(name)
        except FileNotFoundError:
            continue
        except OSError as exc:
            failed.append({"name": name, "error": str(exc)})

return {
    "hostname": gethostname(),
    "pid": os.getpid(),
    "dry_run": dry_run,
    "min_age_days": min_age_days,
    "sem_mp_count": len(items),
    "orphan_old_count": len(orphans),
    "deleted_count": len(deleted),
    "not_found_in_proc_maps": [
        item["name"] for item in items if not holders.get(item["name"])
    ],
    "would_delete": orphans,
    "deleted": deleted,
    "failed": failed,
    "kept_mapped_count": len(kept_mapped),
    "kept_young_count": len(kept_young),
}

预览 1 天前的清理候选项:

Copy
{"min_age_days": 1}

省略 dry_run 时默认为 true,只列出预计删除的条目,不删除文件。完全不传参数时同样默认预览 1 天前的候选项。

核对预览结果并确认可清理后,删除本次扫描中符合条件的文件:

Copy
{"min_age_days": 1, "dry_run": false}

dry_run 必须使用 JSON 布尔值 true 或 false,不要传字符串 "false"、数字或 null。min_age_days 使用非负有限数值,单位为天。删除模式会重新扫描,实际候选项可能与上一次预览不同,并非直接使用上一次的 would_delete 列表。

以下示例中,主机名为示意值,实际以节点返回结果为准:

Copy
{
    "hostname": "deep-pipeline-example",
    "pid": 133,
    "dry_run": true,
    "min_age_days": 1,
    "sem_mp_count": 30,
    "orphan_old_count": 0,
    "deleted_count": 0,
    "not_found_in_proc_maps": [],
    "would_delete": [],
    "deleted": [],
    "failed": [],
    "kept_mapped_count": 30,
    "kept_young_count": 0
}

结果可以按以下方式阅读:

  • hostname、pid:执行节点的主机名和进程 ID,用于核对排查位置。

  • sem_mp_count:扫描到的 sem.mp-* 条目总数。

  • not_found_in_proc_maps:本次未找到进程映射的全部条目,包含修改时间未超过阈值的条目。

  • would_delete、orphan_old_count:未找到映射且修改时间超过阈值的候选列表及数量;字段名中的 orphan 不代表已经确认资源废弃。

  • deleted、deleted_count、failed:实际删除成功的条目、数量及删除失败信息;预览模式下均为空或为 0,failed 不包含扫描时被跳过的进程读取错误。

  • kept_mapped_count:因找到进程映射而保留的数量;kept_young_count:未找到映射,但修改时间未超过阈值而保留的数量。

本例扫描到 30 个条目,全部在可读取的进程映射中找到,因此没有清理候选项,也没有删除任何文件。

Python 环境变量与进程映射的背景说明分别参见 Python os.environ 文档和 Linux /proc/<pid>/maps 手册。

回到顶部

咨询热线

400-821-9199