Repository navigation
Conversation
There was a problem hiding this comment.
这个sender有点重量级了,一些函数可以拆分一下,可以重构为:
transport下有个sender文件夹,然后导出HttpRecordSender,这样测试也好写一些,不过当前PR的目的不是这个,所以可以先写个issue记录,等这个PR合并后再拆分
There was a problem hiding this comment.
对DataStore的改动按我理解是为了避免调用方加一堆If else
不过仓库内既有风格是 Null Object 模式(swanlab/sdk/internal/run/components/null.py 的 NullEmitter/NullConsumer/NullTerminalProxy),writer 的 skip_store 应改用同样方式——新建 NullDataStoreWriter(no-op open/write/skip_records/close),由 core.py 按 skip_store 选择实例,去掉 DataStoreWriter 内部的 skip 标志、write() assert 和 skip_records() 条件分支。
Null Object 模式的好处为避免业务类型内部处理一些非必要 if else,职责也更明确
| """direct-source 模式:不依赖本地镜像目录,直接监听 source_path 所在目录。 | ||
|
|
||
| - 按 source_path.parent 分组,一个目录只 schedule 一次; | ||
| - _registered 以源文件绝对路径为 key,事件精确匹配(同目录其他文件被忽略); |
There was a problem hiding this comment.
这里可以再评估一下为了skip store模式单独开发 direct-source 模式的合理性
主要是会不会引来一些不必要的bug,如果这里不太确定,我的建议是不用添加watcher,和对writer的处理一样新建一个 Null Object
然后在save层warning,只允许now,通过限制行为的方式来减少工作量,并且减少不确定性
There was a problem hiding this comment.
同上一个评论,如果这里很复杂的话也没必要加一些魔法 hhh
| raise TypeError("Object has no len") | ||
|
|
||
|
|
||
| class MemoryViewReader(io.RawIOBase): |
| source_path=ctx.metadata_file.absolute().as_posix(), | ||
| type=SaveType.SAVE_TYPE_METADATA, | ||
| ) | ||
| content = sys_info.metadata.model_dump_json(by_alias=True) |
There was a problem hiding this comment.
probe里面对metadata、requirements、conda有重复的处理逻辑,可以写个helper函数(例如 make_save_record)来解决,这样也能写测试
至于config、builder那边,由于跨模块了,另行处理
| # 也保证 Windows 下文件内容 (LF) 与上方 sha256/size 计算结果一致 | ||
| fs.safe_write(path / filename, content_encode, mode="wb") | ||
| return MediaItem(filename=filename, sha256=sha256, size=len(content_encode), caption=self.caption) | ||
| item = MediaItem( |
There was a problem hiding this comment.
笑死了 感觉可以封装个函数
在 TransformMedia 基类(swanlab/sdk/internal/run/transforms/init.py 或所在基类模块)加一个终态助手:
def _attach_content(self, item: MediaItem, path: Optional[Path], content: bytes) -> MediaItem:
"""skip_store(path=None)时内容进 payload(空内容也赋值以保留 presence),否则落盘。"""
if path is None:
item.payload = content
else:
fs.safe_write(path / item.filename, content, mode="wb")
return item| return self | ||
|
|
||
| @model_validator(mode="after") | ||
| def validate_skip_store(self) -> "Settings": |
Extract a `should_mkdirs` flag to avoid duplicating the mode/skip_store check, and clarify the affected comments and Go formatting in the generated save proto.
|
顺便,测试似乎失败了,看了一下似乎又是windows上的特殊行为,可以记个issue后续修复: 根因分析失败测试:
st = os.stat(path)
return f"{st.st_mtime_ns}:{st.st_size}"测试流程是:写 问题在于:
结果 修复方案(计划)推荐:签名中加入文件标识 # helper.py compute_signature
st = os.stat(path)
return f"{st.st_mtime_ns}:{st.st_size}:{st.st_ino}"
备选(不推荐):只改测试用不同长度的内容(如 验证: |
新增
core.skip_store设置(仅 online 模式合法):开启后 SDK 不在本地产生任何文件(无swanlog/、run-*.swanlab、media/、files/、debug/),全部数据直传云端。Related Issue: #1713
使用
设计
协议:proto 全部追加字段,向后兼容——
MediaItem.payload=5(optional,区分空文件与缺失)、SaveRecord.payload=6、CoreSettings.skip_store=10、ProbeSettings.skip_store=13。本地 DataStore 文件格式版本未变,旧swanlog仍可被新版读取与sync。落盘短路:
DataStoreWriter(skip=True)的 open/write/close 全部无副作用(write 仅计数供 close 统计);init跳过一切目录创建;诊断日志只输出终端。媒体:
transform(path=None)时将内容写入MediaItem.payload而非落盘;sender 经MemoryViewReader零拷贝内存直传对象存储,不触碰本地 media 路径。内部 save(config/metadata/requirements/conda):内容按落盘同款编码内联进
SaveRecord.payload(config 复用dump_config,避免云端结构漂移);sender 解析后走 profile 上传。解析失败属确定性脏数据,告警跳过,不进入 Transport 无限重试;上传失败保持既有 ApiError 分类(5xx 重试 / 4xx 跳过)。CUSTOM save:
payload恒空(非空视为协议违约丢弃),从source_path原路径读取上传,不创建本地软链接镜像;policy="live"的 watcher 直接监听源文件所在目录(direct-source 模式,事件路径精确匹配,同源多 name 一对多注册)。probe:通过新增的显式
ProbeSettings.skip_store字段感知模式,metadata/requirements/conda 直接注入 payload,不写files/目录。模式约束:Settings validator 保证 skip_store 仅 online 合法(含
cloud别名归一化、env 注入、merge 降级路径);交互式引导从 online 降级到 offline 时显式关闭 skip_store 并告警,保证离线数据正常落盘。取舍
swanlab sync/swanlab watch不适用;init 时打印警告明示。测试
单元测试与静态检查
uv run pytest tests/unituv run ruff check .uv run basedpyrightcd core && go build ./...make protoprotos/源同步,无额外 diff覆盖点:
基准测试
新增两个基准(
tests/benchmark/),分别衡量本地持久化层与端到端运行时链路的开销差异。1. 本地持久化层
tests/benchmark/sdk/internal/core_python/store/bench_store_skip.py,直接走生产路径CorePython._store_records,对比落盘(skip_store=False)与完全跳过持久化(skip_store=True)的本地工作:bench_metrics_steps.py对齐),落盘写 LevelDB log,skip 仅计数Image.transform),落盘写media/image/(含 fsync),skip 内联MediaItem.payloadCore._handle_custom_save),落盘建软链接镜像,skip 不建本机(macOS / Apple Silicon,Python 3.11,每类重复 3 次取最优):
落盘侧附加指标:标量吞吐 1.97 M rec/s、媒体写入带宽 140.8 MiB/s、
run-*.swanlab34.4 MB、媒体文件 1,000 个 / 32,768,000 B、save 镜像 100 个;skip 侧无任何本地文件。数据完整性:标量可完整回读,媒体文件数/字节数、save 镜像数均与写入一致。2. 端到端(online + mock HTTP)
tests/benchmark/sdk/cmd/bench_skip_store_e2e.py,完整跑swanlab.init → log/log_image/save → finish,HTTP 全部 mock,对比 skip_store 对用户线程延时与端到端吞吐量的影响。为让 producer / finish 两阶段边界确定,record_interval设为大值,上传统一发生在 finish。负载:1000 step × 100 key = 100,000 标量、250 张媒体、50 个 save。
本机(同上,每场景重复 2 次取最优):
run.log平均延时run.logp50 / p95 / p99解读:
run.log延时基本不变(~113 µs),因为生产端本就只做入队。