Docker中运行AWS IoT Java SDK v2 20分钟后断开,错误码1051
在Docker环境中运行AWS IoT Java SDK v2版本的代码时,运行约20分钟后出现连接断开问题,错误信息为:1051 <[AWS_IO_SOCKET_CLOSED] socket is closed>;但在本地PC或EC2实例上运行该代码时未出现此错误。以下是相关代码及Dockerfile配置,需排查是否为连接/超时配置或Dockerfile存在问题。
代码片段
private final MqttClientConnectionEvents callbacks = new MqttClientConnectionEvents() { public void onConnectionInterrupted(int code) { log.warn("Interruptor <{}> : [{}] {}", code, CRT.awsErrorName(code), CRT.awsErrorString(code)); } public void onConnectionResumed(boolean flag) { log.warn("Resuming: {} [true if the session has been resumed.]", flag); CompletableFuture.runAsync(() -> subscribe()); } }; @Value("${service.description}") private String description; @Value(value = "${aws.iot.endpoint}") private String endpoint; @Value(value = "${aws.iot.clientId}") private String clientId; @Value(value = "${aws.iot.access}") private String access; @Value(value = "${aws.iot.secret}") private String secret; private CompletableFuture<JSONObject> async(MqttAbstract req) { final MqttMessage message = new MqttMessage(req.getTopic(), req.bytes(), QualityOfService.AT_LEAST_ONCE); executor.execute(() -> { log.info("Pub-ing to [{}]: {}", req.getTopic(), new String(message.getPayload())); this.publish(message); }); return registry.waitForMqtt(req.getBody().getPayload().getRequestId()); } private void subscribe(InternalResponses topic) { try { this.connection.subscribe(topic.getTopic(), QualityOfService.AT_LEAST_ONCE, topic.getHandler()).get(); } catch (Exception exc) { log.error("subscribe [{}]: {}", exc.getClass().getSimpleName(), exc.getMessage()); } } private void subscribe() { log.info("Starting subscribing..."); try { this.subscribe(new InternalResponses(gateways, this.registry)); } catch (Exception error) { log.error("onApplicationEvent [{}]: {}", error.getClass().getSimpleName(), error.getMessage()); } } public void publish(MqttMessage message) { try { log.info("Publish to {}", message.getTopic()); executor.execute(() -> connection.publish(message)); } catch (Throwable e) { log.error("publish: {} - {}", e.getClass().getSimpleName(), e.getMessage()); } } @PostConstruct public void construct() { final CredentialsProvider provider = new StaticCredentialsProvider.StaticCredentialsProviderBuilder() .withAccessKeyId(access.getBytes(StandardCharsets.UTF_8)) .withSecretAccessKey(secret.getBytes(StandardCharsets.UTF_8)) .build(); try (AwsIotMqttConnectionBuilder builder = AwsIotMqttConnectionBuilder .newDefaultBuilder() .withEndpoint(endpoint) .withClientId(clientId) .withConnectionEventCallbacks(callbacks) .withProtocolOperationTimeoutMs(MINUTE) .withCleanSession(true) .withWebsockets(true) .withWebsocketSigningRegion("eu-north-1") .withWebsocketCredentialsProvider(provider)) { this.connection = builder.build(); this.connection.connect() .whenCompleteAsync((connectResult, error) -> this.subscribe()); log.info("{} to [{}] as {}", description, endpoint, clientId); } catch (Exception e) { log.error("postConstruct: {} - {}", e.getClass().getSimpleName(), e.getMessage()); } } @PreDestroy public void destroy() { log.info("Shutting down [{}]", description); }
Dockerfile配置
FROM openjdk:11-jdk COPY build/libs/* / COPY app.sh /app.sh RUN ["chmod", "+x", "/app.sh"] ENTRYPOINT ["/app.sh"]
可能的原因与解决方案
1. 缺少连接保活配置(最可能的原因)
Docker默认bridge网络的NAT层通常会主动断开长时间无流量的连接(超时时间可能远短于AWS IoT默认的30分钟闲置断开时间),而本地/EC2环境的网络规则对闲置连接的宽容度更高。
解决方法:在AwsIotMqttConnectionBuilder中添加Mqtt心跳与TCP保活配置,强制定期发送流量维持连接:
.withKeepAliveIntervalSeconds(30) // 每30秒发送一次Mqtt心跳包 .withTcpKeepAlive(true) // 启用TCP层保活 .withTcpKeepAliveIntervalSeconds(20) // TCP保活探测间隔 .withTcpKeepAliveTimeoutSeconds(10) // TCP保活超时时间
2. 协议操作超时设置过短
当前代码中withProtocolOperationTimeoutMs(MINUTE)(假设MINUTE为60000毫秒)是单个Mqtt操作(如发布/订阅)的超时时间,若Docker网络存在轻微波动,可能触发操作超时间接导致连接异常。建议调整为更长时间,例如5分钟:
.withProtocolOperationTimeoutMs(5 * 60 * 1000)
3. Docker网络模式验证
临时将Docker容器的网络模式改为host(启动命令添加--net=host),若问题消失,可确认是bridge网络的NAT超时导致。生产环境不建议使用host模式,此时应优先通过配置保活解决。
4. 重连逻辑补充
当前重连后已重新执行订阅,但需确保subscribe()方法在重连时能正确处理异常,避免因订阅失败导致后续无法接收消息。可在onConnectionResumed中添加异常捕获日志,便于排查重连后的问题。
内容的提问来源于stack exchange,提问作者Yaroslav Lukianenko

