fix_tenant_task_utf8.mjs 7.9 KB

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