import { writeFileSync } from 'fs'; import { dirname, join } from 'path'; import { fileURLToPath } from 'url'; const path = join(dirname(fileURLToPath(import.meta.url)), '..', 'src/main/java/com/fs/app/task/TenantTaskRunner.java'); const content = `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\u914d\u7f6e\u7684\u591a\u79df\u6237\u6267\u884c\u5b9a\u65f6\u4efb\u52a1\uff1a\u4ece yml \u8bfb\u53d6 tenant-ids\uff08\u683c\u5f0f\u79df\u6237ID:\u5206\u79df\u53f7\uff09\uff0c\u9010\u79df\u6237\u5207\u5e93\u540e\u6267\u884c\u903b\u8f91\u3002 */ @Slf4j @Component public class TenantTaskRunner { @Value("\${tenant-ids:}") private String tenantIdsConfig; /** \u4e3a true \u65f6\u6309 tenant-ids \u904d\u5386\u6267\u884c\uff1b\u4e3a false \u65f6\u4e0d\u5207\u5e93\uff0c\u76f4\u63a5\u6267\u884c\u4e00\u6b21\uff08\u79c1\u6709\u5316\u5355\u5e93\u90e8\u7f72\uff09 */ @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, "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 runForConfiguredTenants(String taskName, Consumer action) { runForConfiguredTenantScopes(taskName, scope -> action.accept(scope.getTenantId())); } public void runForConfiguredTenants(String taskName, Runnable action) { runForConfiguredTenants(taskName, tenantId -> action.run()); } public void runForConfiguredTenantScopes(String taskName, Consumer action) { if (!saasTaskEnabled) { log.debug("[SaaS Task] saas.task.enabled=false\uff0c\u4efb\u52a1 {} \u4e0d\u5207\u5e93\u76f4\u63a5\u6267\u884c", taskName); action.accept(null); return; } List scopes = parseConfiguredTenantScopes(); if (scopes.isEmpty()) { log.warn("[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); } } public void runForConfiguredTenantScopes(String taskName, Runnable action) { runForConfiguredTenantScopes(taskName, scope -> action.run()); } private void runScopesSequential(List scopes, String taskName, Consumer action) { for (TenantTaskScope scope : scopes) { runForOneTenant(scope, taskName, tenantId -> action.accept(scope)); } } 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, tenantId -> action.accept(scope)); } 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) { Long tenantId = scope.getTenantId(); try { sopTenantDataSourceAspect.switchTenant(tenantId); log.info("[SaaS Task] \u5b9a\u65f6\u4efb\u52a1\u5207\u6362\u6570\u636e\u6e90 {}, task={}", scope.logTag(), taskName); action.accept(tenantId); } catch (Exception e) { log.error("[SaaS Task] \u79df\u6237\u6267\u884c\u5f02\u5e38 {}, task={}", scope.logTag(), taskName, e); } finally { sopTenantDataSourceAspect.clear(); } } List parseConfiguredTenantIds() { return parseConfiguredTenantScopes().stream() .map(TenantTaskScope::getTenantId) .distinct() .collect(Collectors.toList()); } 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.warn("[SaaS Task] \u5ffd\u7565\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.warn("[SaaS Task] \u5ffd\u7565\u7f3a\u5c11\u5206\u79df\u53f7\u7684\u914d\u7f6e: {}", entry); return null; } return new TenantTaskScope(tenantId, groupNo); } catch (NumberFormatException e) { log.warn("[SaaS Task] \u5ffd\u7565\u975e\u6cd5\u79df\u6237\u914d\u7f6e: {}", entry); return null; } } } `; writeFileSync(path, content, 'utf8'); console.log('written', path);