BackgroundAccessor SPI

BackgroundAccessor 是由三个 SPI 组成的家族,用于把运行在 HTTP 请求之外的 plugin 代码 —— Kafka consumer、定时任务、重试 worker、爬虫 frontier —— 桥接回那些通常依赖 MetaContext(tenant id、user id、语言)的宿主服务。该家族在 Crawler V2 期间引入,遵循单一模式:

Accessor解析对象典型调用方
BackgroundTenantAccessorTenant id 列表 / per-tenant 配置在 tenant 间遍历的 auto-mode worker
BackgroundDataAccessor按 tenantId 限定作用域的动态表 CRUD(每次调用绑定 tenant 上下文)写入动态表的 Kafka consumer
BackgroundConnectorCredentialAccessor按 PID 查找 connector 凭据需要 API key / cookie 的 fetch worker

三者皆为可选 —— plugin 用 @Autowired(required = false) 注入。当宿主未注册实现(旧版平台或仅 OSS 部署)时,plugin 回退到 V1 行为(通常是:不进行 tenant 分发、单租户模式、不做凭据解析)。

为何使用可选注入

该模式让 plugin 能以单个 jar 同时适配未注册实现的宿主与已接入 TenantService 的宿主。桥由 plugin-apiplatform-plugin-api)中的接口声明;实现作为 Spring bean 在宿主中注册。这避免了将 plugin 代码硬耦合到平台内部的 TenantService / DynamicDataService / ApiConnectorService 类。

接口签名

package com.auraboot.framework.plugin.extension;

public interface BackgroundTenantAccessor {
    /** 所有 active tenant id。空列表 → 单租户回退。 */
    List<Long> listActiveTenantIds();
}

public interface BackgroundDataAccessor {
    /**
     * 每个方法都显式接受 tenantId 作为首个参数,并在调用期间把它压入平台的
     * tenant 上下文。无 `withTenantContext` 包裹 —— 直接调用对应的 CRUD 方法即可。
     * 写入动态表时必需 —— 平台据 tenantId 解析 MetaContext。
     */
    Map<String, Object> create(long tenantId, String modelCode, Map<String, Object> data);

    /**
     * 幂等插入。命中 unique 约束时返回 Optional.empty()(而非抛异常),
     * 让 at-least-once consumer(Kafka、重试)可重放同一逻辑记录。
     */
    Optional<Map<String, Object>> tryCreate(long tenantId, String modelCode, Map<String, Object> data);

    /** 按主键读取单条记录;未找到返回 null。 */
    Map<String, Object> getById(long tenantId, String modelCode, String recordId);

    /** 按 field-code → 值的精确匹配(全部 AND 等值)查询。 */
    List<Map<String, Object>> query(long tenantId, String modelCode, Map<String, Object> filters);

    /** 更新已存在记录的字段。 */
    Map<String, Object> update(long tenantId, String modelCode, String recordId, Map<String, Object> data);

    /** 按主键删除记录。 */
    void delete(long tenantId, String modelCode, String recordId);
}

public interface BackgroundConnectorCredentialAccessor {
    /**
     * 按 PID 查找 connector 凭据快照(base URL / auth type / auth-config JSON /
     * default headers)。PID 未知或无匹配 connector 时返回 Optional.empty()。
     */
    Optional<ConnectorCredentials> lookupByPid(String connectorPid);
}

实现示例

某 plugin 中的 Kafka consumer 将抓取到的文档 upsert 到动态表,并按事件携带的 tenant 限定作用域:

package com.acme.crawler.consumer;

import com.auraboot.framework.plugin.extension.BackgroundDataAccessor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

import java.util.Map;

@Component
public class CrawledDocumentUpsertConsumer {

    @Autowired(required = false)
    private BackgroundDataAccessor dataAccessor; // 在 OSS 上为 null

    @KafkaListener(topics = "crawler.document.upsert")
    public void onMessage(CrawledDocumentEvent event) {
        Map<String, Object> row = Map.of(
            "pid", event.getPid(),
            "title", event.getTitle(),
            "content_sha", event.getSha(),
            "content", event.getContent()
        );

        if (dataAccessor != null) {
            // tenantId 作为首个参数显式传入;平台据此绑定 tenant 上下文
            dataAccessor.create(event.getTenantId(), "cr_crawled_document", row);
        } else {
            // V1 回退:单租户模式,由 plugin 自带的写入路径处理
            // ...
        }
    }
}

一个用于 fetcher 的 connector-credential 示例:

import com.auraboot.framework.plugin.extension.BackgroundConnectorCredentialAccessor;
import com.auraboot.framework.plugin.extension.BackgroundConnectorCredentialAccessor.ConnectorCredentials;

@Autowired(required = false)
private BackgroundConnectorCredentialAccessor credentials;

public Map<String, String> headersFor(String connectorPid) {
    if (credentials == null) return Map.of();
    return credentials.lookupByPid(connectorPid)
        .map(ConnectorCredentials::getDefaultHeaders)
        .orElse(Map.of());
}

注册

宿主端:注册一个实现该接口的 Spring @Bean。Plugin 端:除 @Autowired(required = false) 外无需其他。平台默认提供由 TenantServiceMetaContextDynamicDataService)与 ApiConnectorService 支撑的实现(BackgroundTenantAccessorImpl / BackgroundDataAccessorImpl / BackgroundConnectorCredentialAccessorImpl)。

当你发布使用其中任一 accessor 的 plugin 时,请在 plugin 的 README 中记录回退行为,以便运维理解从 OSS 迁移到企业版时会发生什么变化。

常见陷阱

  • 忘记 required = false Plugin 在 OSS 平台上无法启动。
  • 写入动态表时漏传 tenantId 每个 CRUD 方法都以 tenantId 为首个参数,平台据此绑定 tenant 上下文;务必从事件中取出正确的 tenant,不要硬编码。
  • 在 at-least-once consumer 中用 create 而非 tryCreate Kafka / 重试会重放同一条消息;tryCreate 在命中 unique 约束时返回 Optional.empty() 而不抛异常,便于幂等处理。create 则会抛出。
  • 跨 tenant 缓存 accessor 的结果。 每次调用都按传入的 tenantId 限定作用域。请按事件重新解析,不要复用上一条事件的结果。
  • listActiveTenantIds() 的空结果当作错误。 这是有文档约定的回退信号(如 bootstrap 前的全新安装)—— 记录 info 日志,按单租户运行。

相关