fix_tenant_task_runner_utf8.mjs 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209
  1. import { writeFileSync } from 'fs';
  2. import { dirname, join } from 'path';
  3. import { fileURLToPath } from 'url';
  4. const path = join(dirname(fileURLToPath(import.meta.url)), '..', 'src/main/java/com/fs/app/task/TenantTaskRunner.java');
  5. const content = `package com.fs.app.task;
  6. import com.fs.common.utils.StringUtils;
  7. import com.fs.framework.aspectj.SopTenantDataSourceAspect;
  8. import lombok.extern.slf4j.Slf4j;
  9. import org.springframework.beans.factory.annotation.Autowired;
  10. import org.springframework.beans.factory.annotation.Value;
  11. import org.springframework.stereotype.Component;
  12. import javax.annotation.PostConstruct;
  13. import javax.annotation.PreDestroy;
  14. import java.util.ArrayList;
  15. import java.util.Collections;
  16. import java.util.List;
  17. import java.util.concurrent.CountDownLatch;
  18. import java.util.concurrent.ExecutorService;
  19. import java.util.concurrent.LinkedBlockingQueue;
  20. import java.util.concurrent.ThreadPoolExecutor;
  21. import java.util.concurrent.TimeUnit;
  22. import java.util.function.Consumer;
  23. import java.util.stream.Collectors;
  24. /**
  25. * \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
  26. */
  27. @Slf4j
  28. @Component
  29. public class TenantTaskRunner {
  30. @Value("\${tenant-ids:}")
  31. private String tenantIdsConfig;
  32. /** \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 */
  33. @Value("\${saas.task.enabled:true}")
  34. private boolean saasTaskEnabled;
  35. @Value("\${saas.task.parallel:false}")
  36. private boolean parallelEnabled;
  37. @Value("\${saas.task.parallel.threads:4}")
  38. private int parallelThreads;
  39. @Autowired
  40. private SopTenantDataSourceAspect sopTenantDataSourceAspect;
  41. private ExecutorService tenantExecutor;
  42. @PostConstruct
  43. public void initExecutor() {
  44. if (parallelEnabled && parallelThreads > 0) {
  45. tenantExecutor = new ThreadPoolExecutor(
  46. Math.min(4, parallelThreads),
  47. parallelThreads,
  48. 60L, TimeUnit.SECONDS,
  49. new LinkedBlockingQueue<>(512),
  50. r -> {
  51. Thread t = new Thread(r, "ipad-tenant-task-" + System.identityHashCode(r));
  52. t.setDaemon(false);
  53. return t;
  54. },
  55. new ThreadPoolExecutor.CallerRunsPolicy()
  56. );
  57. log.info("[SaaS Task] \u6309\u79df\u6237\u5e76\u884c\u5df2\u5f00\u542f\uff0c\u7ebf\u7a0b\u6c60\u5927\u5c0f: {}", parallelThreads);
  58. }
  59. }
  60. @PreDestroy
  61. public void shutdownExecutor() {
  62. if (tenantExecutor != null) {
  63. tenantExecutor.shutdown();
  64. try {
  65. if (!tenantExecutor.awaitTermination(30, TimeUnit.SECONDS)) {
  66. tenantExecutor.shutdownNow();
  67. }
  68. } catch (InterruptedException e) {
  69. tenantExecutor.shutdownNow();
  70. Thread.currentThread().interrupt();
  71. }
  72. }
  73. }
  74. public void runForConfiguredTenants(String taskName, Consumer<Long> action) {
  75. runForConfiguredTenantScopes(taskName, scope -> action.accept(scope.getTenantId()));
  76. }
  77. public void runForConfiguredTenants(String taskName, Runnable action) {
  78. runForConfiguredTenants(taskName, tenantId -> action.run());
  79. }
  80. public void runForConfiguredTenantScopes(String taskName, Consumer<TenantTaskScope> action) {
  81. if (!saasTaskEnabled) {
  82. log.debug("[SaaS Task] saas.task.enabled=false\uff0c\u4efb\u52a1 {} \u4e0d\u5207\u5e93\u76f4\u63a5\u6267\u884c", taskName);
  83. action.accept(null);
  84. return;
  85. }
  86. List<TenantTaskScope> scopes = parseConfiguredTenantScopes();
  87. if (scopes.isEmpty()) {
  88. log.warn("[SaaS Task] \u672a\u914d\u7f6e tenant-ids\uff0c\u8df3\u8fc7\u4efb\u52a1 {}", taskName);
  89. return;
  90. }
  91. if (parallelEnabled && tenantExecutor != null) {
  92. runScopesParallel(scopes, taskName, action);
  93. } else {
  94. runScopesSequential(scopes, taskName, action);
  95. }
  96. }
  97. public void runForConfiguredTenantScopes(String taskName, Runnable action) {
  98. runForConfiguredTenantScopes(taskName, scope -> action.run());
  99. }
  100. private void runScopesSequential(List<TenantTaskScope> scopes, String taskName, Consumer<TenantTaskScope> action) {
  101. for (TenantTaskScope scope : scopes) {
  102. runForOneTenant(scope, taskName, tenantId -> action.accept(scope));
  103. }
  104. }
  105. private void runScopesParallel(List<TenantTaskScope> scopes, String taskName, Consumer<TenantTaskScope> action) {
  106. CountDownLatch latch = new CountDownLatch(scopes.size());
  107. for (TenantTaskScope scope : scopes) {
  108. tenantExecutor.submit(() -> {
  109. try {
  110. runForOneTenant(scope, taskName, tenantId -> action.accept(scope));
  111. } finally {
  112. latch.countDown();
  113. }
  114. });
  115. }
  116. try {
  117. if (!latch.await(30, TimeUnit.MINUTES)) {
  118. log.warn("[SaaS Task] \u6309\u79df\u6237\u5e76\u884c\u6267\u884c\u8d85\u65f6(30\u5206\u949f), task={}", taskName);
  119. }
  120. } catch (InterruptedException e) {
  121. Thread.currentThread().interrupt();
  122. log.warn("[SaaS Task] \u6309\u79df\u6237\u5e76\u884c\u6267\u884c\u88ab\u4e2d\u65ad, task={}", taskName, e);
  123. }
  124. }
  125. private void runForOneTenant(TenantTaskScope scope, String taskName, Consumer<Long> action) {
  126. Long tenantId = scope.getTenantId();
  127. try {
  128. sopTenantDataSourceAspect.switchTenant(tenantId);
  129. log.info("[SaaS Task] \u5b9a\u65f6\u4efb\u52a1\u5207\u6362\u6570\u636e\u6e90 {}, task={}", scope.logTag(), taskName);
  130. action.accept(tenantId);
  131. } catch (Exception e) {
  132. log.error("[SaaS Task] \u79df\u6237\u6267\u884c\u5f02\u5e38 {}, task={}", scope.logTag(), taskName, e);
  133. } finally {
  134. sopTenantDataSourceAspect.clear();
  135. }
  136. }
  137. List<Long> parseConfiguredTenantIds() {
  138. return parseConfiguredTenantScopes().stream()
  139. .map(TenantTaskScope::getTenantId)
  140. .distinct()
  141. .collect(Collectors.toList());
  142. }
  143. public List<TenantTaskScope> parseConfiguredTenantScopes() {
  144. if (StringUtils.isBlank(tenantIdsConfig)) {
  145. return Collections.emptyList();
  146. }
  147. List<TenantTaskScope> scopes = new ArrayList<>();
  148. for (String part : tenantIdsConfig.split(",")) {
  149. if (StringUtils.isBlank(part)) {
  150. continue;
  151. }
  152. TenantTaskScope scope = parseTenantScopeEntry(part.trim());
  153. if (scope != null) {
  154. scopes.add(scope);
  155. }
  156. }
  157. return new ArrayList<>(scopes.stream()
  158. .collect(Collectors.toMap(
  159. scope -> scope.getTenantId() + ":" + scope.getGroupNo(),
  160. scope -> scope,
  161. (left, right) -> left))
  162. .values());
  163. }
  164. private TenantTaskScope parseTenantScopeEntry(String entry) {
  165. int colonIndex = entry.indexOf(':');
  166. if (colonIndex <= 0) {
  167. log.warn("[SaaS Task] \u5ffd\u7565\u975e\u6cd5\u914d\u7f6e\uff08\u5fc5\u987b\u4e3a \u79df\u6237ID:\u5206\u79df\u53f7\uff09: {}", entry);
  168. return null;
  169. }
  170. try {
  171. Long tenantId = Long.valueOf(entry.substring(0, colonIndex).trim());
  172. String groupNo = entry.substring(colonIndex + 1).trim();
  173. if (StringUtils.isBlank(groupNo)) {
  174. log.warn("[SaaS Task] \u5ffd\u7565\u7f3a\u5c11\u5206\u79df\u53f7\u7684\u914d\u7f6e: {}", entry);
  175. return null;
  176. }
  177. return new TenantTaskScope(tenantId, groupNo);
  178. } catch (NumberFormatException e) {
  179. log.warn("[SaaS Task] \u5ffd\u7565\u975e\u6cd5\u79df\u6237\u914d\u7f6e: {}", entry);
  180. return null;
  181. }
  182. }
  183. }
  184. `;
  185. writeFileSync(path, content, 'utf8');
  186. console.log('written', path);