Telegram Bot生产环境无法向十万级用户推送消息求助
问题分析与排查方案
结合你的生产环境场景(10万+消息量)和代码实现,我推测核心问题大概率和Telegram API的速率限制、响应结果的不完整处理有关,下面分点拆解并给出优化建议:
1. 优先排查Telegram API的速率限制
测试环境400用户量小,不会触发Telegram Bot API的速率阈值,但生产环境10万+消息量肯定会触发限制。你的代码目前只处理了200和403(用户拉黑)状态,但完全忽略了429(请求过多)这个高频场景:
- Telegram返回
429时会附带retry_after字段,告诉你需要等待多久才能继续请求 - 你的代码会把
429判定为发送失败,但没有做延迟重试,导致消息被反复重发,最终超过MAX_REDELIVERY后被强制确认,用户自然收不到
验证方式:立刻在analizeResponse方法中添加全量响应日志,把所有Telegram返回的内容都打印出来,比如在读取content后加上:
log.info("Full Telegram Response: {}", content.toString());
如果日志里出现429和retry_after字段,就坐实了速率限制问题。
2. 完善响应结果的处理逻辑
你的analizeResponse只返回boolean,无法区分「需要重试的失败」和「永久失败」,这在高并发场景下非常致命。建议重构这个方法,返回更丰富的结果:
重构后的响应处理代码
// 定义辅助类承载响应结果 class TelegramResponse { private boolean isSuccess; private boolean needRetry; private int retryDelaySeconds; // 仅当needRetry为true时有效 public TelegramResponse(boolean isSuccess, boolean needRetry, int retryDelaySeconds) { this.isSuccess = isSuccess; this.needRetry = needRetry; this.retryDelaySeconds = retryDelaySeconds; } // Getter方法 public boolean isSuccess() { return isSuccess; } public boolean isNeedRetry() { return needRetry; } public int getRetryDelaySeconds() { return retryDelaySeconds; } } private TelegramResponse analyzeResponse(List<NameValuePair> payload, HttpResponse response) throws IOException { HttpEntity entity = response.getEntity(); int statusCode = response.getStatusLine().getStatusCode(); String responseContent = ""; if (entity != null) { BufferedReader in = new BufferedReader(new InputStreamReader(entity.getContent(), "UTF-8")); StringBuilder content = new StringBuilder(); String inputLine; while ((inputLine = in.readLine()) != null) { content.append(inputLine); } in.close(); responseContent = content.toString(); log.info("Telegram API Response - Status: {}, Content: {}", statusCode, responseContent); } // 处理200状态:注意Telegram返回200时也可能包含ok:false的错误 if (statusCode == 200) { ObjectMapper mapper = new ObjectMapper(); JsonNode root = mapper.readTree(responseContent); boolean apiOk = root.get("ok").asBoolean(); if (apiOk) { return new TelegramResponse(true, false, 0); } else { String errorDesc = root.get("description").asText(); log.error("Telegram API returned 200 but failed: {}", errorDesc); // 处理用户拉黑的情况 if (errorDesc.contains("bot was blocked by the user")) { String chatId = payload.stream() .filter(p -> p.getName().equals("chat_id")) .findFirst() .map(NameValuePair::getValue) .orElse(null); if (chatId != null) { new TelegramUserDao().unsubscribeUser(chatId); } } // 判断是否需要重试:临时错误重试,永久错误直接放弃 boolean needRetry = !errorDesc.contains("bot was blocked") && !errorDesc.contains("invalid chat ID"); return new TelegramResponse(false, needRetry, needRetry ? 5 : 0); } } // 处理403:用户拉黑或权限问题,直接放弃 else if (statusCode == 403) { if (responseContent.contains("Forbidden: bot was blocked by the user")) { String chatId = payload.stream() .filter(p -> p.getName().equals("chat_id")) .findFirst() .map(NameValuePair::getValue) .orElse(null); if (chatId != null) { new TelegramUserDao().unsubscribeUser(chatId); } } return new TelegramResponse(false, false, 0); } // 处理429:速率限制,按Telegram要求延迟重试 else if (statusCode == 429) { ObjectMapper mapper = new ObjectMapper(); JsonNode root = mapper.readTree(responseContent); int retryAfter = root.get("retry_after").asInt(); log.warn("Rate limited by Telegram, retry after {} seconds", retryAfter); return new TelegramResponse(false, true, retryAfter); } // 其他错误:默认重试5秒 else { log.error("Unexpected Telegram API status code: {}", statusCode); return new TelegramResponse(false, true, 5); } }
消费者逻辑适配
修改ActiveMQ消费者的处理逻辑,根据返回结果决定是否重试:
TelegramResponse response = telegramSender.sendMessage(telegramBody); if (response.isSuccess()) { sent++; message.acknowledge(); } else { noSent++; int deliveryCount = message.getIntProperty(DELIVERY_COUNT); // 如果需要重试且未达到最大重试次数,延迟后不确认消息(让ActiveMQ重新入队) if (response.isNeedRetry() && deliveryCount < MAX_REDELIVERY) { Thread.sleep(response.getRetryDelaySeconds() * 1000); // 注意:这里不要调用message.acknowledge(),让消息重新进入队列等待下一次消费 } else { message.acknowledge(); log.warn("News {} failed to send to user {} after {} attempts", news.getId(), newsLetter.getUserId(), deliveryCount); } }
3. 其他潜在排查点
- Chat ID有效性:生产环境可能存在大量无效Chat ID(用户注销、未启动过对话等),建议定期批量校验Chat ID的有效性,清理无效数据减少无效请求。
- 网络差异:确认生产环境和测试环境的网络是否一致,比如是否有防火墙、代理限制了Telegram API的访问。
- ActiveMQ配置:检查ActiveMQ的重发策略,比如是否设置了合理的重发间隔,避免短时间内重复发送触发更严格的速率限制。
内容的提问来源于stack exchange,提问作者Анатолий Житарь
相关产品推荐
相关产品推荐

