| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209 |
- 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<Long> 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<TenantTaskScope> 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<TenantTaskScope> 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<TenantTaskScope> scopes, String taskName, Consumer<TenantTaskScope> action) {
- for (TenantTaskScope scope : scopes) {
- runForOneTenant(scope, taskName, tenantId -> action.accept(scope));
- }
- }
- 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, 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<Long> 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<Long> parseConfiguredTenantIds() {
- return parseConfiguredTenantScopes().stream()
- .map(TenantTaskScope::getTenantId)
- .distinct()
- .collect(Collectors.toList());
- }
- 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.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);
|