用于执行输入的 PY 代码,输出结果为非数据集。
输入 Python 代码即可。
提供联想功能并支持变量引用,具体语法参见预置变量,使用 Tab 键接受联想:
若多个PY代码节点都用到相同的方法,推荐在全局设置-PY设置中配置公共脚本,具体配置参见全局设置。
PY 代码适合处理参数、状态和小型结果,不建议承载大批量明细数据。大数据量应尽量保留为数据集,并使用查询、数据转换或 SQL 转换节点处理,避免转换为 Python 列表、字典或 Pandas DataFrame 后造成额外内存占用。更多建议参见并发、队列与内存优化。
案例:多个 PY 节点都使用DataTableMySQL方法。
在全局设置-PY设置中配置公共脚本:
PY代码节点可以直接调用,无需再重复定义:
数据流运行引擎基于 Python,PY代码节点除了处理业务参数,也可以调用 Python 标准库查看运行环境、辅助排查问题。可执行的操作取决于节点实际运行环境中安装的库、可访问的目录和进程权限。
在数据流中添加一个PY代码节点,在执行代码中输入:
import os
dict(os.environ)
执行节点后,返回结果就是当前节点 Python 进程可见的全部环境变量,以字典形式展示。这里查看的是节点实际运行所在容器或进程的环境,不是浏览器所在电脑的环境,也不代表其他 Pod 或整个集群的环境变量。
如果只需确认某个变量,可以按名称查询,例如:
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 读取本次数据流的启动参数,参数配置参见预置变量。
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 天前的清理候选项:
{"min_age_days": 1}
省略 dry_run 时默认为 true,只列出预计删除的条目,不删除文件。完全不传参数时同样默认预览 1 天前的候选项。
核对预览结果并确认可清理后,删除本次扫描中符合条件的文件:
{"min_age_days": 1, "dry_run": false}
dry_run 必须使用 JSON 布尔值 true 或 false,不要传字符串 "false"、数字或 null。min_age_days 使用非负有限数值,单位为天。删除模式会重新扫描,实际候选项可能与上一次预览不同,并非直接使用上一次的 would_delete 列表。
以下示例中,主机名为示意值,实际以节点返回结果为准:
{
"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 手册。
回到顶部
咨询热线
