OrderReminderTask.java
13.4 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
package com.infoloop.tianting.logic.task;
import com.infoloop.tianting.deliveryruleservice.Rule;
import com.infoloop.tianting.mealorderservice.PendingOrderSubscriptionMessageRpcResponse;
import com.infoloop.tianting.mealorderservice.SingleMpAccountRpcResponse;
import com.infoloop.tianting.model.dto.WxUserDto.SendSubscriptionMessageResponse;
import com.infoloop.tianting.service.OrderReminderService;
import com.infoloop.tianting.service.client.DeliveryRuleServiceRpcClient;
import com.infoloop.tianting.service.client.MealOrderServiceRpcClient;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections4.CollectionUtils;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.time.DayOfWeek;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executor;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;
import static com.infoloop.tianting.constant.ConfigConstants.ASYNC_EXECUTOR;
/**
* 订餐提醒定时任务
* 在订餐周期开启前(周一 9:00)向已授权的家长发送微信服务通知
*
* 推送逻辑:
* 1. 查询所有待发送的订阅消息(未删除的记录)
* 2. 按 openId 去重,每个用户只推送一次
* 3. 通过 openId 查询活跃的小程序用户
* 4. 根据配送规则配置生成订餐时间提示
* 5. 发送消息,发送成功后标记为已删除
*/
@Slf4j
@Component
public class OrderReminderTask {
private static final int ENTERPRISE_ID = 449;
/**
* 温馨提醒内容
*/
private static final String TIPS = "请点击卡片,准时进入小程序完成预订";
/**
* 订餐时间配置
*/
private record OrderTimeConfig(String startTime, String endTime) {}
private final MealOrderServiceRpcClient mealOrderServiceRpcClient;
private final DeliveryRuleServiceRpcClient deliveryRuleServiceRpcClient;
private final OrderReminderService orderReminderService;
private final Executor asyncExecutor;
public OrderReminderTask(
MealOrderServiceRpcClient mealOrderServiceRpcClient,
DeliveryRuleServiceRpcClient deliveryRuleServiceRpcClient,
OrderReminderService orderReminderService,
@Qualifier(ASYNC_EXECUTOR) Executor asyncExecutor) {
this.mealOrderServiceRpcClient = mealOrderServiceRpcClient;
this.deliveryRuleServiceRpcClient = deliveryRuleServiceRpcClient;
this.orderReminderService = orderReminderService;
this.asyncExecutor = asyncExecutor;
}
/**
* 每周一 9:00 执行订餐提醒推送
*/
@Scheduled(cron = "0 0 9 * * MON")
public void executeTask() {
log.info("OrderReminderTask started");
try {
// 1. 获取配送规则配置,用于生成订餐时间提示
OrderTimeConfig orderTimeConfig = getOrderTimeConfigFromRule();
if (orderTimeConfig == null) {
log.warn("Failed to get order time config from rule, using default");
orderTimeConfig = new OrderTimeConfig("请查看小程序", "请查看小程序");
}
log.info("Order time config: startTime={}, endTime={}", orderTimeConfig.startTime(), orderTimeConfig.endTime());
// 2. 查询所有待发送的订阅消息(未删除的记录)
List<PendingOrderSubscriptionMessageRpcResponse> pendingMessages =
mealOrderServiceRpcClient.queryPendingOrderSubscriptionMessages(ENTERPRISE_ID);
if (pendingMessages.isEmpty()) {
log.info("No pending order subscription messages to send");
return;
}
log.info("Found {} pending order subscription messages", pendingMessages.size());
// 3. 提取所有 openId,查询活跃的小程序用户
Set<String> openIds = pendingMessages.stream()
.map(PendingOrderSubscriptionMessageRpcResponse::getOpenId)
.collect(Collectors.toSet());
Map<String, SingleMpAccountRpcResponse> activeMpAccountMap = getActiveMpAccountMap(openIds);
log.info("Found {} active mp accounts from {} openIds", activeMpAccountMap.size(), openIds.size());
int totalMessages = pendingMessages.size();
AtomicInteger skippedInactive = new AtomicInteger(0);
AtomicInteger skippedDuplicate = new AtomicInteger(0);
AtomicInteger sentCount = new AtomicInteger(0);
List<Long> sentMessageIds = Collections.synchronizedList(new ArrayList<>());
Set<String> sentOpenIds = ConcurrentHashMap.newKeySet(); // 线程安全的 Set,用于 openId 去重
// 4. 按 openId 去重,筛选出需要发送的消息
List<PendingOrderSubscriptionMessageRpcResponse> messagesToSend = new ArrayList<>();
for (PendingOrderSubscriptionMessageRpcResponse message : pendingMessages) {
String openId = message.getOpenId();
// 检查用户是否活跃
if (!activeMpAccountMap.containsKey(openId)) {
log.debug("Skipping message {} - user {} is not active", message.getId(), openId);
skippedInactive.incrementAndGet();
sentMessageIds.add(message.getId());
continue;
}
// openId 去重:每个用户只推送一次
if (!sentOpenIds.add(openId)) {
log.debug("Skipping message {} - already queued for openId: {}", message.getId(), openId);
skippedDuplicate.incrementAndGet();
sentMessageIds.add(message.getId());
continue;
}
messagesToSend.add(message);
}
log.info("Messages to send: {}, Skipped inactive: {}, Skipped duplicate: {}",
messagesToSend.size(), skippedInactive.get(), skippedDuplicate.get());
// 5. 多线程并发发送消息
final OrderTimeConfig finalOrderTimeConfig = orderTimeConfig;
List<CompletableFuture<Void>> futures = messagesToSend.stream()
.map(message -> CompletableFuture.runAsync(() -> {
try {
boolean success = sendReminderMessage(message.getOpenId(), message, finalOrderTimeConfig);
if (success) {
sentMessageIds.add(message.getId());
sentCount.incrementAndGet();
}
} catch (Exception e) {
log.error("Failed to send reminder to openId: {}, messageId: {}",
message.getOpenId(), message.getId(), e);
}
}, asyncExecutor))
.toList();
// 等待所有任务完成
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
// 6. 批量标记已发送的消息(逻辑删除)
if (!sentMessageIds.isEmpty()) {
int affectedRows = mealOrderServiceRpcClient.batchDeleteOrderSubscriptionMessagesByIds(
ENTERPRISE_ID, sentMessageIds);
log.info("Marked {} messages as sent", affectedRows);
}
log.info("OrderReminderTask completed. Total: {}, Sent: {}, SkippedInactive: {}, SkippedDuplicate: {}",
totalMessages, sentCount.get(), skippedInactive.get(), skippedDuplicate.get());
} catch (Exception e) {
log.error("OrderReminderTask failed", e);
}
}
/**
* 从配送规则配置中获取订餐时间配置
* 返回开始时间和结束时间,格式:周x xx:xx
*/
private OrderTimeConfig getOrderTimeConfigFromRule() {
try {
final var response = deliveryRuleServiceRpcClient.getDefaultDeliveryTemplateAndRule(ENTERPRISE_ID);
final var rule = response.getRule();
if (rule == null || rule.getConfigJson().getRulesList().isEmpty()) {
log.warn("Delivery rule not configured");
return null;
}
final var cfg = rule.getConfigJson().getRules(0);
final String startTimeStr = cfg.getStartTime();
final String endTimeStr = cfg.getEndTime();
if (startTimeStr == null || endTimeStr == null) {
log.warn("Delivery rule time config incomplete");
return null;
}
// 获取启用的日期
List<DayOfWeek> enabledDays = getEnabledDays(cfg);
if (enabledDays.isEmpty()) {
log.warn("No enabled days in delivery rule");
return null;
}
// 格式化时间(去掉秒)
String startTime = formatTime(startTimeStr);
String endTime = formatTime(endTimeStr);
// 获取第一个和最后一个启用日
DayOfWeek firstDay = enabledDays.get(0);
DayOfWeek lastDay = enabledDays.get(enabledDays.size() - 1);
// 计算本周的实际日期(微信time类型需要具体日期格式)
LocalDate today = LocalDate.now();
LocalDate startDate = today.with(java.time.temporal.TemporalAdjusters.nextOrSame(firstDay));
LocalDate endDate = today.with(java.time.temporal.TemporalAdjusters.nextOrSame(lastDay));
// 如果结束日期在开始日期之前,说明跨周了,结束日期加一周
if (endDate.isBefore(startDate)) {
endDate = endDate.plusWeeks(1);
}
// 格式化为微信接受的格式:yyyy年MM月dd日 HH:mm
java.time.format.DateTimeFormatter dateFormatter = java.time.format.DateTimeFormatter.ofPattern("yyyy年MM月dd日");
String formattedStartTime = startDate.format(dateFormatter) + " " + startTime;
String formattedEndTime = endDate.format(dateFormatter) + " " + endTime;
return new OrderTimeConfig(formattedStartTime, formattedEndTime);
} catch (Exception e) {
log.error("Failed to get order time config from rule", e);
return null;
}
}
/**
* 获取启用的日期列表(按周一到周日排序)
*/
private List<DayOfWeek> getEnabledDays(Rule cfg) {
List<DayOfWeek> enabledDays = new ArrayList<>();
if (Boolean.TRUE.equals(cfg.getMonday())) enabledDays.add(DayOfWeek.MONDAY);
if (Boolean.TRUE.equals(cfg.getTuesday())) enabledDays.add(DayOfWeek.TUESDAY);
if (Boolean.TRUE.equals(cfg.getWednesday())) enabledDays.add(DayOfWeek.WEDNESDAY);
if (Boolean.TRUE.equals(cfg.getThursday())) enabledDays.add(DayOfWeek.THURSDAY);
if (Boolean.TRUE.equals(cfg.getFriday())) enabledDays.add(DayOfWeek.FRIDAY);
if (Boolean.TRUE.equals(cfg.getSaturday())) enabledDays.add(DayOfWeek.SATURDAY);
if (Boolean.TRUE.equals(cfg.getSunday())) enabledDays.add(DayOfWeek.SUNDAY);
return enabledDays;
}
/**
* 格式化时间字符串(HH:mm:ss -> HH:mm)
*/
private String formatTime(String timeStr) {
if (timeStr == null || timeStr.length() < 5) {
return timeStr;
}
// 取前5位 HH:mm
return timeStr.substring(0, 5);
}
/**
* 获取活跃的小程序用户 Map(openId -> MpAccount)
*/
private Map<String, SingleMpAccountRpcResponse> getActiveMpAccountMap(Set<String> openIds) {
if (CollectionUtils.isEmpty(openIds)) {
return Collections.emptyMap();
}
List<SingleMpAccountRpcResponse> allAccounts = mealOrderServiceRpcClient.queryAllActiveMpAccounts(ENTERPRISE_ID);
return allAccounts.stream()
.filter(account -> openIds.contains(account.getOpenId()))
.collect(Collectors.toMap(SingleMpAccountRpcResponse::getOpenId, account -> account, (a, b) -> a));
}
/**
* 发送单条提醒消息
*
* 小程序模板字段:
* - time4: 开始时间(周x xx:xx)
* - time5: 结束时间(周x xx:xx)
* - thing1: 温馨提醒
*/
private boolean sendReminderMessage(String openId, PendingOrderSubscriptionMessageRpcResponse message, OrderTimeConfig orderTimeConfig) {
String jumpPath = message.getJumpPath();
// 发送消息
SendSubscriptionMessageResponse response = orderReminderService.sendOrderReminder(
openId,
jumpPath,
orderTimeConfig.startTime(), // time4: 开始时间
orderTimeConfig.endTime(), // time5: 结束时间
TIPS // thing1: 温馨提醒
);
// 判断发送结果
if (response.getErrcode() != null && response.getErrcode() == 0) {
return true;
}
// 记录失败原因
log.warn("Failed to send reminder. openId: {}, errcode: {}, errmsg: {}",
openId, response.getErrcode(), response.getErrmsg());
return false;
}
}