| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217 |
- 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<TenantTaskScope> 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<TenantTaskScope> 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<TenantTaskScope> scopes, String taskName, Consumer<TenantTaskScope> action) {
- for (TenantTaskScope scope : scopes) {
- runForOneTenant(scope, taskName, action);
- }
- }
- private void runScopesParallel(List<TenantTaskScope> scopes, String taskName, Consumer<TenantTaskScope> 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<TenantTaskScope> 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<TenantTaskScope> parseConfiguredTenantScopes() {
- if (StringUtils.isBlank(tenantIdsConfig)) {
- return Collections.emptyList();
- }
- List<TenantTaskScope> 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');
|