BISHENG 频道信息源订阅状态治理:同步订阅 + 每日对账(F031 实战解析)

BISHENG 频道信息源订阅状态治理:同步订阅 + 每日对账(F031 实战解析) BISHENG 频道信息源订阅状态治理同步订阅 每日对账F031 实战解析【免费下载链接】bishengBISHENG is an open LLM devops platform for next generation Enterprise AI applications. Powerful and comprehensive features include: GenAI workflow, RAG, Agent, Unified model management, Evaluation, SFT, Dataset Management, Enterprise-level System Management, Observability and more.项目地址: https://gitcode.com/GitHub_Trending/bi/bisheng本文导读BISHENG 开放 LLM DevOps 平台仓库根目录 README_CN.md的频道Channel模块需要与外部情报服务bisheng_information_client保持信息源订阅一致。v2.6.0 的 F031 特性针对频道模块暴露的三类缺陷——重复订阅、共用信息源被过度退订导致文章静默断更、channel_info_source悬挂行——给出了一套同步订阅 每日对账的最终一致性方案。读完本文你将掌握该特性从验收标准、架构决策、Repository/Service 改造、Celery Beat 薄包装到 Test-First 测试矩阵的完整实现脉络并能直接对照源码文件在仓库中定位每一处实现。本文主体基于特性文档 tasks.md 与配套规格 spec.md并辅以仓库源码核验实现事实。1. 背景与问题定义F031 源自 v2.6.0-beta3 频道模块的缺陷收敛没有独立 PRD优先级为 P1所属版本 v2.6.0合入feat/2.6.0-beta3分支。其要解决的三个核心问题是重复订阅新建频道选择已被其他频道订阅过的信息源时系统仍会重复发起订阅请求既浪费情报服务调用也让订阅上限19007校验对已订阅源重复计数过度退订删除dismiss或编辑update频道时同步执行退订导致仍被其他频道使用的共用信息源被一起退订造成频道文章静默断更悬挂行同步退订后channel_info_source表可能残留不再被引用、也不再订阅的信息源行表无界增长。与之对应的是三条用户故事spec §1频道使用者/创建者删除或编辑其中一个频道时不应把仍被其它频道使用的信息源一起退订避免文章静默断更频道创建者新建频道时系统不再对已订阅的源重复发起订阅请求19007 上限校验只对真正新增的源计数后端/运维订阅判定走索引查询而非JSON_CONTAINS全表扫退订与悬挂行清理收敛到每日凌晨一次批量对账。1.1 与用户对齐的三条前提Spec Discovery 阶段与用户对齐的关键不确定性结论退订改为每日对账驱动退订延迟 ≤ 1 天可接受对账粒度先做纯每日全量暂不做事件触发 每日兜底的近实时方案订阅仍保持同步即时——19007 上限必须在建频道前校验且用户期望建完很快有文章。1.2 范围边界本次纳入spec §0确立「订阅意图唯一真相 租户内所有channel.source_list的并集」channel_info_source定位为「已订阅集合 展示元数据」的物化视图行存在 ⟺ 毕昇认定该源已订阅订阅判据切到channel_info_source主键索引查询取代find_channels_by_source_id的 JSON 全表扫去掉dismiss_channel/update_channel移除源时的同步退订新增每日对账Celery Beat 任务desired ↔ channel_info_source ↔ 外部情报服务三方收敛退订 兜底补订阅 删悬挂行ChannelInfoSourceRepository新增delete_by_ids。明确排除事件触发近实时对账、主动查询外部订阅态做反向校正情报服务无查询接口、前端改动、find_channels_by_source_id方法的删除保留但不再用于订阅判据。2. 验收标准AC 全表AC-ID 在本特性内唯一格式AC-NNtasks.md 中的测试任务通过覆盖 AC: AC-NN追溯到此表。这是理解后续任务拆解与测试用例的总纲。2.1 同步订阅热路径ID角色操作预期结果AC-01已登录用户create_channelsource_list中部分源在channel_info_source已存在、部分不存在仅对不存在的源调用subscribe_information_source参数恰为缺失源集合保持入参顺序去重已存在的源不调用AC-02已登录用户create_channel全部源在channel_info_source已存在不调用subscribe_information_source正常建频道AC-03已登录用户create_channel对新源subscribe抛InformationSourceSubscriptionLimitError(19007)在持久化频道之前中止不产生 channel / membership / OpenFGA owner 元组AC-04已登录用户update_channelto_add中部分源已在channel_info_source仅订阅to_add中缺失的源已存在的不订阅AC-05系统订阅成功后对每个新订阅的源向channel_info_source插入元数据行id / name / icon / type / description主键冲突时幂等跳过2.2 退订改由对账驱动热路径不退订ID角色操作预期结果AC-06已登录用户dismiss_channel频道含source_list不调用unsubscribe_information_source不删除channel_info_source行仅删频道与关系AC-07已登录用户update_channelto_remove非空不调用unsubscribe_information_source移除仅改channel.source_listAC-08共用源场景频道 A、B 都引用源 Xdismiss AX 仍订阅、channel_info_source中 X 行仍在B 文章不受影响2.3 每日对账最终一致性兜底ID角色操作预期结果AC-09对账任务某源 X 不在任何channel.source_listcurrent - desired调用unsubscribe_information_source([X])且删除channel_info_source中 X 行AC-10对账任务某源 Y 在desired但不在channel_info_source漏订阅 / 同步失败残留调用subscribe_information_source([Y])、补元数据、插入行AC-11对账任务desired current不产生任何情报服务调用、不增删行AC-12对账任务多租户多个活跃租户按租户分别对账desired/current/外部 API-key 均在各自租户上下文内收敛无跨租户串源AC-13对账任务单源unsubscribe或subscribe抛异常记录日志logger.exception不中断其余源的对账逐源/分批隔离失败本轮失败项留待下一轮自愈3. 边界情况与已知盲区spec §3 明确列出以下边界处理策略并发建频道引用同一新源两请求都查到channel_info_source无该源 → 都subscribe外部按集合幂等→ 都尝试插行主键冲突时一方幂等跳过。属可接受的 best-effort一致性由对账兜底不引入分布式锁同步订阅成功、插行失败进程崩溃等外部已订阅但本地无行 → 下次对账desired - current命中 → 幂等重订阅 补行自愈退订延迟源不再被引用后最长到下一次对账才退订≤ 1 天其间多拉文章无害。不支持「立即强制退订」如有合规需求需另设手动触发本期延后外部情报服务静默丢订阅对账只比desired vs channel_info_source不查外部无法发现此类漂移情报服务无查询接口。缓解方案是对desired全集做幂等重订阅可选开关默认关避免每天全量重订阅压力本期默认不开记为已知盲区source_list含重复 id订阅判据按dict.fromkeys去重后处理不支持信息源元数据name/icon的实时刷新——channel_info_source行存在即跳过重拉元数据可能轻微滞后可接受。4. 架构决策AD 全表ID决策选项结论理由AD-01订阅意图的唯一真相A:channel.source_list并集 / B: 独立计数列 / C:channel_info_sourceA为真相channel_info_source为其物化视图并集即真相不引入需双写维护的计数列避免漂移AD-02订阅判据数据源A:find_channels_by_source_id(JSON_CONTAINS) / B:channel_info_source.find_by_ids(主键索引)BA 不可索引、N 次全表扫违反项目 dual-DB「禁 JSON_CONTAINS」规则B 为 O(1) 索引命中AD-03订阅 vs 退订时效A: 二者皆同步 / B: 订阅同步、退订对账B不对称订阅需即时19007 校验 文章时效退订不紧急交对账可同时根治「过度退订」AD-04channel_info_source行生命周期A: 创建即插、退订即删 / B: 与外部订阅态同生共死均由对账在 1→0 时删B让「行存在 ⟺ 外部已订阅」恒成立消除「立刻退订却留行 → 漏订阅」的旧反例AD-05对账触发方式A: 纯每日全量 / B: 事件触发 每日兜底A本期满足最终一致诉求最简B 延后AD-06对账失败隔离A: 整批事务 / B: 逐源或分批隔离 下轮自愈B单源外部调用失败不应阻断其余源对账本就幂等可重入4.1 关键不变量AD-04 的同生共死AD-04 是整个方案正确性的基石。旧语义下立刻退订却留行会造成漏订阅反例dismiss_channel同步调用unsubscribe但channel_info_source行未删后续create_channel用find_by_ids判据发现行存在就跳过订阅——于是外部已退订、本地却认为已订阅形成永久漏订阅。AD-04 改为行与外部订阅态同生共死均由对账在 1→0 时删后行存在 ⟺ 外部已订阅恒成立该反例被消除。5. 数据库与 Domain 模型不新增表、不改channel_info_source表结构。现有模型见 channel_info_source.py含tenant_id多租户自动隔离id 信息源 id 主键。channel.source_listJSON 列维持现状仍作为每个频道引用源的列表与文章过滤依据。仅在 Repository 接口层增加删除能力spec §5、T001# src/backend/bisheng/channel/domain/repositories/interfaces/channel_info_source_repository.py abstractmethod async def delete_by_ids(self, source_ids: List[str]) - None: Delete metadata rows for the given source ids (reconcile cleanup).从当前接口文件 channel_info_source_repository.py 可以看到ChannelInfoSourceRepository最终具备五个核心能力find_by_ids索引判据、batch_add补行、get_by_page分页取用、find_all对账全量取 currentT006 预留实现、delete_by_ids对账清理。5.1 接口无新表、无迁移本特性无 DDL / 无表结构变更因此无需迁移与回滚多租户隔离依赖仓库自动注入两张表均含tenant_id调用方保证已处于目标租户上下文即可。6. API 契约无新增对外 API对账为Celery Beat 周期任务非 HTTP 入口无新增 / 修改的对外 API退订与对账均为内部行为沿用错误码InformationSourceSubscriptionLimitError(19007)不新增错误码、无需在 release-contract.md 注册新模块编码模块编码复用 channel 模块 190 段Owner F026。7. Service 层逻辑与调用时机对账业务逻辑落在domain servicechannel_service.pyCelery 任务仅作薄包装做租户遍历 调用遵守worker → service → repo分层便于脱离 Celery 直接单测。7.1 核心方法一览方法位置输入输出职责create_channel改channel_service.pyCreateDTO UserPayloadChannel用channel_info_source.find_by_ids求缺失源 → 仅订阅缺失源持久化前→ 持久化 → 补元数据行update_channel改channel_service.pyUpdateDTO UserPayloadChannelto_add仅订阅缺失源to_remove不退订仅改 source_list补元数据行dismiss_channel改channel_service.pychannel_id UserPayloadbool删频道与关系移除对unsubscribe_information_source的调用reconcile_information_subscriptions新servicechannel_service.py当前租户上下文对账统计to_sub/to_unsub/failed计算 desired并集(source_list)、currentchannel_info_sourcecurrent-desired退订删行desired-current订阅补行逐源隔离失败reconcile_all_tenants新Celery 薄包装reconcile.py——Beat 入口遍历活跃租户、逐租户注入上下文后调用上面的 service 方法不含业务逻辑7.2 调用时机总表落地后入口subscribeunsubscribecreate_channel仅缺失源持久化前索引判据—update_channel 新增源仅缺失源—update_channel 移除源——交对账dismiss_channel——交对账成员退订 / 移除成员——本就不涉及信息源每日对账desired - current兜底补订阅current - desired退订7.3 源码级实现核验create_channel 的索引判据channel_service.py先find_by_ids(channel_data.source_list)取已订阅集合再用dict.fromkeys去重求missing_source_ids仅当存在缺失源时才subscribe_information_source(missing_source_ids)并补元数据——且订阅发生在ChannelRepository.save之前保证 19007 上限错误中止时不留下孤儿 channel / membership / OpenFGA owner 元组AC-03。update_channel 的差分逻辑channel_service.py以集合差计算to_add_sources new - old与to_remove_sources old - newto_add再经find_by_ids过滤出真正缺失的missing_add才订阅注释明确说明Removed sources are NOT unsubscribed here——移除仅体现在channel.source_list更新与source_list_changed标记后者触发update_channels_latest_article_time退订交每日对账AC-07。而更新后的元数据同步走_sync_channel_info_source_metadata对新增源还会调度sync_information_article.apply_async延后 1 小时拉取文章。dismiss_channel删除末尾的unsubscribe_information_source(channel.source_list)调用块。可以在全文件搜索验证unsubscribe_information_source在 channel_service.py 中仅剩reconcile_information_subscriptions内部一处调用约 L438证明 dismiss/update 热路径已彻底不再退订AC-06 / AC-08。7.4 对账方法实现T006 落地reconcile_information_subscriptions的实现channel_service.py与 spec §7.2 伪流程完全一致desired await self.channel_repository.find_all_referenced_source_ids() # T002: source_list 并集 current_rows await self.channel_info_source_repository.find_all() # T006 偏差1: 一次全量 current {row.id for row in current_rows} to_unsub current - desired # 孤儿退订 删行 to_sub desired - current # 缺失订阅 补元数据 for source_id in to_unsub: try: await bisheng_information_client.unsubscribe_information_source([source_id]) await self.channel_info_source_repository.delete_by_ids([source_id]) except Exception: logger.exception(...) # 逐源失败隔离 failed 1 for source_id in to_sub: try: await bisheng_information_client.subscribe_information_source([source_id]) await self._sync_channel_info_source_metadata(bisheng_information_client, [source_id]) except Exception: logger.exception(...) failed 1 return {to_sub: len(to_sub), to_unsub: len(to_unsub), failed: failed}find_all_referenced_source_ids的实现channel_repository_impl.py不用JSON_CONTAINS/JSON_EXTRACT——只SELECT Channel.source_list将 JSON 列原样取出Python 内展开求并集兼顾 DM8 兼容全表读仅发生在每日一次的对账中而非每次写操作。_sync_channel_info_source_metadatachannel_service.pyT004 抽出的私有 helper负责拉元数据id/name/icon/business_type/description并batch_add被 create/update/reconcile 三处复用避免重复update 的文章定时同步sync_information_article.apply_async仍保留在 update_channel 内、不进 helper偏差 2。7.5 权限 / 多租户无新增资源授权不涉及PermissionService.authorize()channel/channel_info_source均多租户自动隔离对账由 Beat 逐租户注入上下文参考既有 article.py 的租户处理与默认celery队列约定。8. WorkerCelery Beat 薄包装8.1 任务与调度注册新增文件 reconcile.py 定义reconcile_all_tenants任务注册方式参考worker/information/article.py与 worker/config.pybeat_schedule settings.celery_task.beat_schedule调度来自 Celery 配置。任务体不含业务逻辑仅租户遍历 调用 service。任务执行的关键细节T008 约束走默认celery队列审批/文章同步同款不配task_routesBeat 凌晨低峰执行定在04:3030 4 * * *偏差 4按评审建议与文章同步 Beat05:30错开避免叠加多租户放大坑单租户对账失败不影响其余租户外层再包一层 try/except 日志。8.2 逐租户上下文注入_reconcile_one_tenant用set_current_tenant_id(tenant_id)ContextVarBeat 内部迭代注入非 Celery headers 传参切换租户上下文活跃租户集合由_active_tenant_ids()给出——多租户关闭时回退到默认租户DEFAULT_TENANT_ID即单租户部署等价 tenant_id1启用时取ROOT_TENANT_ID 活跃子租户列表。worker 构造 service 时通过_channel_service_session()异步上下文管理器开get_async_db_session其中space_channel_member_repositoryNone对账方法不依赖它——这是 spec 未细化的实现细节偏差 3。9. 任务拆解Tasks与测试矩阵开发模式为后端 Test-First先写测试红→ 再写实现绿。本特性纯后端、无新表、无对外 API、无前端、无新错误码。9.1 闭环约束重要本特性的三件事——①订阅判据切到channel_info_source、②去掉 dismiss/update 的同步退订、③上线每日对账——必须同批合入。单独换判据会在旧退订语义下漏订阅详见 spec §4 AD-04 / §9 注。当前已合入的过渡实现create_channel用find_channels_by_source_id将在 T003 被替换。9.2 基础设施无测试配对T001ChannelInfoSourceRepository.delete_by_ids——按id IN (...)删除元数据行对账to_unsub清理用空列表直接返回不发 SQL实现见 channel_info_source_repository_impl.py。T002ChannelRepository.find_all_referenced_source_ids——读取当前租户所有channel.source_list列Python 内展开求并集返回去重集合对账的desired全表读仅在每日对账批量调用一次可接受。9.3 同步订阅判据 去同步退订T003单测红新建 test_channel_source_subscription.py并修改 test_channel_relation_compat.py。Mockbisheng_information_client、channel_info_source_repository、channel_repository用例与 AC 映射用例断言覆盖 ACtest_create_subscribes_only_missing_sourcessource_list[A,B]、find_by_ids返回[A]→ subscribe 仅以[B]调用一次AC-01test_create_skips_subscribe_when_all_present全部已在表内 → subscribe 不被调用AC-02test_create_aborts_before_persist_on_limitsubscribe 抛 19007 → channel/membership/owner 元组均未写AC-03test_create_inserts_metadata_rows_for_new_sourcesget_information_source_by_idsbatch_add以缺失源调用AC-05test_update_add_subscribes_only_missingto_add部分已存在 → 仅订阅缺失AC-04test_update_remove_does_not_unsubscribeto_remove非空 → unsubscribe 不被调用AC-07test_dismiss_does_not_unsubscribedismiss 含 source_list 频道 → unsubscribe 与delete_by_ids均不被调用AC-06test_dismiss_shared_source_stays_subscribedA、B 都引用 Xdismiss A → 整个流程完全不调 unsubscribe、不删 X 行AC-08基础设施复用test_channel_relation_compat.py既有_service()/_LoginUser模式SimpleNamespaceAsyncMock无需新 conftest。注意 AC-03 用例是从既有test_create_channel_aborts_before_persist_when_subscription_limit_exceeded迁移而来仅把channel_repositorymock 从find_channels_by_source_id改为channel_info_source.find_by_ids语义。T004绿改造 channel_service.py 中create_channel/update_channel/dismiss_channel三个方法的信息源订阅副作用见 §7.3 实现核验。跨 Feature 提示本文件与 F026-channel-active-authorization 共享但二者领域解耦F026 拥有频道授权/space_channel_memberchannel 字段写行为F031 拥有channel_info_source订阅生命周期已在 release-contract 表 1 登记本任务不碰授权相关代码。find_channels_by_source_id方法保留不删本特性仅不再在订阅判据中使用它。9.4 每日对账T005单测红新建 test_information_subscription_reconcile.py直接测ChannelService.reconcile_information_subscriptions()mock 三个依赖不经 Celery用例断言覆盖 ACtest_reconcile_unsubscribes_and_deletes_orphandesired{A}、current{A,X} → unsubscribe([X]) delete_by_ids([X])A 不动AC-09test_reconcile_subscribes_and_inserts_missingdesired{A,Y}、current{A} → subscribe([Y]) 拉元数据 batch_add 插 Y 行AC-10test_reconcile_noop_when_equaldesiredcurrent → 无任何调用/增删行AC-11test_reconcile_isolates_per_source_failure单源 unsubscribe 抛异常 → 其余源照常、方法不抛、返回统计含 failed 且记日志AC-13test_reconcile_returns_counts返回{to_sub, to_unsub, failed}统计结构支撑可观测§10T006绿实现reconcile_information_subscriptions见 §7.4复用 T004 抽出的私有 helper 共用「订阅 补元数据」内部逻辑。9.5 Beat 薄包装T007单测红新建 test_information_reconcile_worker.pymock 活跃租户来源 service 方法用例断言覆盖 ACtest_reconcile_all_tenants_iterates_each_tenant3 个活跃租户 → service 各调用一次且每次调用前current_tenant_id被设为对应租户AC-12test_reconcile_all_tenants_isolates_tenant_failure某租户对账抛异常 → 其余租户照常、任务不整体失败、记日志AC-12T008绿定义reconcile_all_tenantsCelery Beat 任务 调度注册见 §8。10. 非功能要求与可观测性性能热路径订阅判定为主键索引查询取代 JSON 全表扫source_list并集的全表扫描每日仅在对账中执行一次非每写操作 N 次正确性消除「共用源被过度退订」与「已订阅源被重复订阅」订阅态每日自愈兼容性不改表结构、不改对外 API、不改前端沿用现有错误码与 Celery 默认队列可观测对账每轮输出to_sub/to_unsub计数与失败明细日志logger.exception便于排查情报服务调用异常多租户Beat 逐租户对账注意「任务数 × 租户数」的既有放大坑安排在凌晨低峰单次执行。11. 实际偏差记录与 Code Review 修订tasks.md 记录了实现相对 spec 的偏差与评审修订是理解真实落地细节的关键补充偏差 1新增 repo 方法除delete_by_ids外给ChannelInfoSourceRepository同时加了find_all()spec T006 已预留对账用它一次性取全量 current偏差 2共享 helperT004 抽出ChannelService._sync_channel_info_source_metadata()拉元数据 batch_add被 create/update/reconcile 三处复用update 的文章定时同步仍保留在 update_channel 内、不进 helper偏差 3worker 构造 service 的方式worker 用_channel_service_session()异步上下文管理器开get_async_db_session并构造ChannelService其中space_channel_member_repositoryNone偏差 4对账时间定在 04:3030 4 * * *早于文章同步 05:30未做spec 已声明排除事件触发近实时对账、外部订阅态反向校正、对账对 desired 全量幂等重订阅默认关均未实现符合 spec §3 范围边界code-review 修订2026-06-04删除死代码ChannelRepository.find_channels_by_source_id接口 impl及 impl 中随之失效的json_array_contains导入——本特性后已无任何调用方ChannelInfoSourceRepository.batch_add改为捕获IntegrityError后回滚 重查去重 仅插新行落实 spec §3 / AC-05 声称的主键冲突幂等跳过并发同时新增同一信息源场景新增test_batch_add_idempotent_on_integrity_error覆盖回滚/重试逻辑。batch_add的幂等实现可从 channel_info_source_repository_impl.py 直接验证session.add_allcommit捕获IntegrityError→rollback→ 重查已存在 id → 仅插remaining从而让并发场景下的插入幂等跳过而非报错中断调用方。12. 总结三层一致性的设计骨架F031 最终落成的机制可以概括为三层协作同步热路径订阅create_channel/update_channel以channel_info_source主键索引为判据仅对缺失源订阅19007 上限在持久化前校验延迟冷路径退订 兜底每日 04:30 的 Celery Beat 任务逐租户执行reconcile_information_subscriptionsdesiredchannel.source_list 并集与currentchannel_info_source 全量求差孤儿退订 删行、缺失补订阅 插行逐源失败隔离、下轮自愈物化视图一致性channel_info_source行与外部订阅态同生共死行存在 ⟺ 外部已订阅配合并发batch_add的幂等跳过消除了重复订阅、过度退订与悬挂行三类缺陷。整个特性以 Test-First 方式推进8/8 任务完成全部测试通过13 条 AC 均被对应单测覆盖测试文件集中在 src/backend/test/channel 目录。相关实现文件索引领域服务channel_service.py仓储接口channel_info_source_repository.py、channel_repository.py仓储实现channel_info_source_repository_impl.py、channel_repository_impl.pyBeat 任务reconcile.py单测test_channel_source_subscription.py、test_information_subscription_reconcile.py、test_information_reconcile_worker.py【免费下载链接】bishengBISHENG is an open LLM devops platform for next generation Enterprise AI applications. Powerful and comprehensive features include: GenAI workflow, RAG, Agent, Unified model management, Evaluation, SFT, Dataset Management, Enterprise-level System Management, Observability and more.项目地址: https://gitcode.com/GitHub_Trending/bi/bisheng创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考