Repository navigation
Conversation
There was a problem hiding this comment.
这个sender有点重量级了,一些函数可以拆分一下,可以重构为:
transport下有个sender文件夹,然后导出HttpRecordSender,这样测试也好写一些,不过当前PR的目的不是这个,所以可以先写个issue记录,等这个PR合并后再拆分
| """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): |
| 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}"
备选(不推荐):只改测试用不同长度的内容(如 验证: |
On Windows/NTFS, a file deleted and recreated with same-size content can keep an identical mtime_ns (timestamp granularity) and os.stat may serve the stale directory-entry cache, so _process_change judged the new file as unchanged and skipped the callback. Including st_ino (NTFS file index) makes delete-then-recreate always yield a different signature.
Refactor media transforms to share a common `_attach_content` helper on `TransformMedia`. This unifies the `skip_store` and file-writing logic across audio, image, video, HTML, text, and other media types, ensuring content is stored in `payload` when no path is provided and written to disk otherwise while preserving empty-content presence semantics.
新增
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),因为生产端本就只做入队。