import { writeFileSync } from 'node:fs'; import { join, dirname } from 'node:path'; import { fileURLToPath } from 'node:url'; const root = join(dirname(fileURLToPath(import.meta.url)), '..'); const runner = `package com.fs.app.task; import com.fs.common.utils.StringUtils; import com.fs.framework.aspectj.SopTenantDataSourceAspect; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.function.Consumer; import java.util.stream.Collectors; /** * \u6309 tenant-ids \u914d\u7f6e\u6267\u884c\u591a\u79df\u6237\u5b9a\u65f6\u4efb\u52a1\uff08\u683c\u5f0f\uff1a\u79df\u6237ID:\u5206\u79df\u53f7\uff09\u3002 */ @Slf4j @Component public class TenantTaskRunner { @Value("\${tenant-ids:}") private String tenantIdsConfig; @Value("\${saas.task.enabled:true}") private boolean saasTaskEnabled; @Value("\${saas.task.parallel:false}") private boolean parallelEnabled; @Value("\${saas.task.parallel.threads:4}") private int parallelThreads; @Autowired private SopTenantDataSourceAspect sopTenantDataSourceAspect; private ExecutorService tenantExecutor; @PostConstruct public void initExecutor() { if (parallelEnabled && parallelThreads > 0) { tenantExecutor = new ThreadPoolExecutor( Math.min(4, parallelThreads), parallelThreads, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(512), r -> { Thread t = new Thread(r, "wx-ipad-tenant-task-" + System.identityHashCode(r)); t.setDaemon(false); return t; }, new ThreadPoolExecutor.CallerRunsPolicy() ); log.info("[SaaS Task] \u6309\u79df\u6237\u5e76\u884c\u5df2\u5f00\u542f\uff0c\u7ebf\u7a0b\u6c60\u5927\u5c0f: {}", parallelThreads); } } @PreDestroy public void shutdownExecutor() { if (tenantExecutor != null) { tenantExecutor.shutdown(); try { if (!tenantExecutor.awaitTermination(30, TimeUnit.SECONDS)) { tenantExecutor.shutdownNow(); } } catch (InterruptedException e) { tenantExecutor.shutdownNow(); Thread.currentThread().interrupt(); } } } public void runForConfiguredTenantScopes(String taskName, Consumer action) { if (!saasTaskEnabled) { log.error("[SaaS Task] saas.task.enabled=false\uff0c\u4efb\u52a1 {} \u672a\u6267\u884c\uff08\u8bf7\u914d\u7f6e tenant-ids \u5e76\u5f00\u542f saas.task.enabled\uff09", taskName); return; } List scopes = parseConfiguredTenantScopes(); if (scopes.isEmpty()) { log.error("[SaaS Task] \u672a\u914d\u7f6e tenant-ids\uff0c\u8df3\u8fc7\u4efb\u52a1 {}", taskName); return; } if (parallelEnabled && tenantExecutor != null) { runScopesParallel(scopes, taskName, action); } else { runScopesSequential(scopes, taskName, action); } } private void runScopesSequential(List scopes, String taskName, Consumer action) { for (TenantTaskScope scope : scopes) { runForOneTenant(scope, taskName, action); } } private void runScopesParallel(List scopes, String taskName, Consumer action) { CountDownLatch latch = new CountDownLatch(scopes.size()); for (TenantTaskScope scope : scopes) { tenantExecutor.submit(() -> { try { runForOneTenant(scope, taskName, action); } finally { latch.countDown(); } }); } try { if (!latch.await(30, TimeUnit.MINUTES)) { log.warn("[SaaS Task] \u6309\u79df\u6237\u5e76\u884c\u6267\u884c\u8d85\u65f6(30\u5206\u949f), task={}", taskName); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); log.warn("[SaaS Task] \u6309\u79df\u6237\u5e76\u884c\u6267\u884c\u88ab\u4e2d\u65ad, task={}", taskName, e); } } private void runForOneTenant(TenantTaskScope scope, String taskName, Consumer action) { try { sopTenantDataSourceAspect.switchTenant(scope.getTenantId()); log.info("[SaaS Task] \u5207\u6362\u6570\u636e\u6e90 {}, task={}", scope.logTag(), taskName); action.accept(scope); } catch (Exception e) { log.error("[SaaS Task] \u79df\u6237\u6267\u884c\u5f02\u5e38 {}, task={}", scope.logTag(), taskName, e); } finally { sopTenantDataSourceAspect.clear(); } } public List parseConfiguredTenantScopes() { if (StringUtils.isBlank(tenantIdsConfig)) { return Collections.emptyList(); } List scopes = new ArrayList<>(); for (String part : tenantIdsConfig.split(",")) { if (StringUtils.isBlank(part)) { continue; } TenantTaskScope scope = parseTenantScopeEntry(part.trim()); if (scope != null) { scopes.add(scope); } } return new ArrayList<>(scopes.stream() .collect(Collectors.toMap( scope -> scope.getTenantId() + ":" + scope.getGroupNo(), scope -> scope, (left, right) -> left)) .values()); } private TenantTaskScope parseTenantScopeEntry(String entry) { int colonIndex = entry.indexOf(':'); if (colonIndex <= 0) { log.error("[SaaS Task] \u975e\u6cd5\u914d\u7f6e\uff08\u5fc5\u987b\u4e3a \u79df\u6237ID:\u5206\u79df\u53f7\uff09: {}", entry); return null; } try { Long tenantId = Long.valueOf(entry.substring(0, colonIndex).trim()); String groupNo = entry.substring(colonIndex + 1).trim(); if (StringUtils.isBlank(groupNo)) { log.error("[SaaS Task] \u7f3a\u5c11\u5206\u79df\u53f7: {}", entry); return null; } return new TenantTaskScope(tenantId, groupNo); } catch (NumberFormatException e) { log.error("[SaaS Task] \u975e\u6cd5\u79df\u6237\u914d\u7f6e: {}", entry); return null; } } } `; const scope = `package com.fs.app.task; import lombok.AllArgsConstructor; import lombok.Getter; import lombok.ToString; /** * \u5b9a\u65f6\u4efb\u52a1\u6267\u884c\u8303\u56f4\uff1a\u7531 tenant-ids \u89e3\u6790\u5f97\u5230\u7684\u79df\u6237 ID \u4e0e\u5206\u79df\u53f7\u3002 */ @Getter @ToString @AllArgsConstructor public class TenantTaskScope { private final Long tenantId; private final String groupNo; public String logTag() { return "tenantId=" + tenantId + ", groupNo=" + groupNo; } public static String logTag(TenantTaskScope scope) { return scope == null ? "tenantId=null, groupNo=null" : scope.logTag(); } } `; writeFileSync(join(root, 'src/main/java/com/fs/app/task/TenantTaskRunner.java'), runner, 'utf8'); writeFileSync(join(root, 'src/main/java/com/fs/app/task/TenantTaskScope.java'), scope, 'utf8'); console.log('done');