You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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,提问作者Анатолий Житарь

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 09:03:39