帮你快速理解、总结文档立即下载

MQTT over QUIC

最近更新时间:2026-09-29 12:00:30
我的收藏

背景

MQTT 协议仅需要传输层提供有序、无损、双向的通信。
说明:
The MQTT protocol requires an underlying transport that provides an ordered, lossless, stream of bytes from the Client to Server and Server to Client.
而 QUIC 协议能够满足 MQTT 对传输层的要求,作为 HTTP 3.0 的传输层协议,它有以下关键特性:
更低延迟:QUIC 设计融合了 TLS,实现了 0-RTT 等特性,有更低的连接延迟。
多路复用:QUIC 基于 UDP 实现,多个 stream 之间没有 TCP 的 head-of-line-blocking 问题。
连接迁移:QUIC 使用 Connection ID 维护会话,当五元组发生变更,可以通过连接迁移恢复会话,避免业务层的有状态协议会话如 MQTT Session 会话重建。
弱网支持:更优秀的丢包处理算法、拥塞算法,让 QUIC 更适合弱网环境使用。

MQTT over QUIC

消息队列 MQTT 版已深度集成并完整支持了 QUIC 协议的核心特性,旨在为您的业务提供更快速、更稳定、更灵活的连接体验。

使用限制

1. 仅支持 IETF RFC9000 QUIC,不支持早期的非标准版本(如 gQUIC)
2. 当前仅支持单流(Single-stream)模式,即 MQTT 的控制 Packet 与数据 Packet 在同一个流中传输,暂未开放使用独立多流进行分离传输的模式。
3. MQTT over QUIC 现已对 1.4.0 及以上版本的集群实例开放。如需升级您的实例版本以体验该特性,欢迎联系我们。

接入点

协议
端口
ALPN
QUIC
14567
mqtt
公网端口
VPC 端口
2. 在左侧导航栏选择资源管理 > 集群管理,选择好地域后,单击目标集群的 ID,进入集群基本信息页面。
3. 如果您已开启公网,则可以在接入信息模块查看当前公网 QUIC 接入地址。

2. 在左侧导航栏选择资源管理 > 集群管理,选择好地域后,单击目标集群的 ID,进入集群基本信息页面。
3. 在接入信息模块查看具体的私有网络 VPC 连接列表,单击目标私有连接,即可展开该私有网络下的 QUIC 接入地址。


SDK

类型
Repository
QUIC 支持
C
C++
Java
Netty 4.2
Android
Netty 4.2

C 示例代码

Publish
Subscribe
Fallback
/*******************************************************************************
* Copyright (c) 2023-2025 EMQ Technologies Co., William Yang and others.
*
* All rights reserved. This program and the accompanying materials
* are made available under the terms of the Eclipse Public License v2.0
* and Eclipse Distribution License v1.0 which accompany this distribution.
*
* The Eclipse Public License is available at
* https://www.eclipse.org/legal/epl-2.0/
* and the Eclipse Distribution License is available at
* http://www.eclipse.org/org/documents/edl-v10.php.
*/

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include "MQTTAsync.h"

#if !defined(_WIN32)
#include <unistd.h>
#else
#include <windows.h>
#endif

#if defined(_WRS_KERNEL)
#include <OsWrapper.h>
#endif

#define ADDRESS "quic://127.0.0.1:14567"
#define CLIENTID "ExampleClientPub"
#define TOPIC "MQTT Examples"
#define PAYLOAD "Hello World!"
#define QOS 1
#define TIMEOUT 10000L

int finished = 0;

/* kept at file scope so the reconnect in connlost() can reuse them */
static const char* g_username = NULL;
static const char* g_password = NULL;
static MQTTAsync_SSLOptions g_sslopts = MQTTAsync_SSLOptions_initializer;

void connlost(void *context, char *cause)
{
MQTTAsync client = (MQTTAsync)context;
MQTTAsync_connectOptions conn_opts = MQTTAsync_connectOptions_initializer;
int rc;

printf("\\nConnection lost\\n");
printf(" cause: %s\\n", cause);

printf("Reconnecting\\n");
conn_opts.keepAliveInterval = 20;
conn_opts.cleansession = 1;
conn_opts.username = g_username;
conn_opts.password = g_password;
conn_opts.ssl = &g_sslopts;
conn_opts.context = client;
if ((rc = MQTTAsync_connect(client, &conn_opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start connect, return code %d\\n", rc);
finished = 1;
}
}

void onDisconnectFailure(void* context, MQTTAsync_failureData* response)
{
printf("Disconnect failed\\n");
finished = 1;
}

void onDisconnect(void* context, MQTTAsync_successData* response)
{
printf("Successful disconnection\\n");
finished = 1;
}

void onSendFailure(void* context, MQTTAsync_failureData* response)
{
MQTTAsync client = (MQTTAsync)context;
MQTTAsync_disconnectOptions opts = MQTTAsync_disconnectOptions_initializer;
int rc;

printf("Message send failed token %d error code %d\\n", response->token, response->code);
opts.onSuccess = onDisconnect;
opts.onFailure = onDisconnectFailure;
opts.context = client;
if ((rc = MQTTAsync_disconnect(client, &opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start disconnect, return code %d\\n", rc);
exit(EXIT_FAILURE);
}
}

void onSend(void* context, MQTTAsync_successData* response)
{
MQTTAsync client = (MQTTAsync)context;
MQTTAsync_disconnectOptions opts = MQTTAsync_disconnectOptions_initializer;
int rc;

printf("Message with token value %d delivery confirmed\\n", response->token);
opts.onSuccess = onDisconnect;
opts.onFailure = onDisconnectFailure;
opts.context = client;
if ((rc = MQTTAsync_disconnect(client, &opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start disconnect, return code %d\\n", rc);
exit(EXIT_FAILURE);
}
}


void onConnectFailure(void* context, MQTTAsync_failureData* response)
{
printf("Connect failed, rc %d\\n", response ? response->code : 0);
finished = 1;
}


void onConnect(void* context, MQTTAsync_successData* response)
{
MQTTAsync client = (MQTTAsync)context;
MQTTAsync_responseOptions opts = MQTTAsync_responseOptions_initializer;
MQTTAsync_message pubmsg = MQTTAsync_message_initializer;
int rc;

printf("Successful connection\\n");
opts.onSuccess = onSend;
opts.onFailure = onSendFailure;
opts.context = client;
pubmsg.payload = PAYLOAD;
pubmsg.payloadlen = (int)strlen(PAYLOAD);
pubmsg.qos = QOS;
pubmsg.retained = 0;
if ((rc = MQTTAsync_sendMessage(client, TOPIC, &pubmsg, &opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start sendMessage, return code %d\\n", rc);
exit(EXIT_FAILURE);
}
}

int messageArrived(void* context, char* topicName, int topicLen, MQTTAsync_message* m)
{
/* not expecting any messages */
return 1;
}

void trace_callback(enum MQTTASYNC_TRACE_LEVELS level, char* message)
{
fprintf(stdout, "Trace : %d, %s\\n", level, message);
}

int main(int argc, char* argv[])
{
MQTTAsync client;
MQTTAsync_setTraceCallback(trace_callback);
MQTTAsync_setTraceLevel(MQTTASYNC_TRACE_MINIMUM);
//MQTTAsync_setTraceLevel(MQTTASYNC_TRACE_PROTOCOL);
MQTTAsync_connectOptions conn_opts = MQTTAsync_connectOptions_initializer;
g_sslopts.enableServerCertAuth = 0; //for simplicity, we don't verify the server certificate
int rc;

const char* uri = (argc > 1) ? argv[1] : ADDRESS;
g_username = (argc > 2) ? argv[2] : NULL;
g_password = (argc > 3) ? argv[3] : NULL;


if ((rc = MQTTAsync_create(&client, uri, CLIENTID, MQTTCLIENT_PERSISTENCE_NONE, NULL)) != MQTTASYNC_SUCCESS)
{
printf("Failed to create client object, return code %d\\n", rc);
exit(EXIT_FAILURE);
}

if ((rc = MQTTAsync_setCallbacks(client, NULL, connlost, messageArrived, NULL)) != MQTTASYNC_SUCCESS)
{
printf("Failed to set callback, return code %d\\n", rc);
exit(EXIT_FAILURE);
}

conn_opts.keepAliveInterval = 20;
conn_opts.cleansession = 1;
conn_opts.username = g_username;
conn_opts.password = g_password;
conn_opts.onSuccess = onConnect;
conn_opts.onFailure = onConnectFailure;
conn_opts.httpProxy = NULL;
conn_opts.httpsProxy = NULL;
conn_opts.context = client;
conn_opts.ssl = &g_sslopts;
if ((rc = MQTTAsync_connect(client, &conn_opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start connect, return code %d\\n", rc);
exit(EXIT_FAILURE);
}

printf("Waiting for publication of %s\\n"
"on topic %s for client with ClientID: %s\\n",
PAYLOAD, TOPIC, CLIENTID);
while (!finished)
#if defined(_WIN32)
Sleep(100);
#else
usleep(10000L);
#endif

printf("Destroy client\\n");
MQTTAsync_destroy(&client);
return rc;
}

/*******************************************************************************
* Copyright (c) 2023-2025 EMQ Technologies Co., William Yang and others.
*
* All rights reserved. This program and the accompanying materials
* are made available under the terms of the Eclipse Public License v2.0
* and Eclipse Distribution License v1.0 which accompany this distribution.
*
* The Eclipse Public License is available at
* https://www.eclipse.org/legal/epl-2.0/
* and the Eclipse Distribution License is available at
* http://www.eclipse.org/org/documents/edl-v10.php.
*/

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include "MQTTAsync.h"

#if !defined(_WIN32)
#include <unistd.h>
#else
#include <windows.h>
#endif

#if defined(_WRS_KERNEL)
#include <OsWrapper.h>
#endif

#define ADDRESS "quic://127.0.0.1:14567"
#define CLIENTID "ExampleClientSub"
#define TOPIC "MQTT Examples"
#define QOS 1
#define TIMEOUT 10000L

int disc_finished = 0;
int subscribed = 0;
int finished = 0;

/* kept at file scope so the reconnect in connlost() can reuse them */
static const char* g_username = NULL;
static const char* g_password = NULL;
static MQTTAsync_SSLOptions g_sslopts = MQTTAsync_SSLOptions_initializer;

void onConnect(void* context, MQTTAsync_successData* response);
void onConnectFailure(void* context, MQTTAsync_failureData* response);

void connlost(void *context, char *cause)
{
MQTTAsync client = (MQTTAsync)context;
MQTTAsync_connectOptions conn_opts = MQTTAsync_connectOptions_initializer;
int rc;

printf("\\nConnection lost\\n");
if (cause)
printf(" cause: %s\\n", cause);

printf("Reconnecting\\n");
conn_opts.keepAliveInterval = 20;
conn_opts.cleansession = 1;
conn_opts.username = g_username;
conn_opts.password = g_password;
conn_opts.ssl = &g_sslopts;
conn_opts.onSuccess = onConnect;
conn_opts.onFailure = onConnectFailure;
conn_opts.context = client;
if ((rc = MQTTAsync_connect(client, &conn_opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start connect, return code %d\\n", rc);
finished = 1;
}
}


int msgarrvd(void *context, char *topicName, int topicLen, MQTTAsync_message *message)
{
printf("Message arrived\\n");
printf(" topic: %s\\n", topicName);
printf(" message: %.*s\\n", message->payloadlen, (char*)message->payload);
MQTTAsync_freeMessage(&message);
MQTTAsync_free(topicName);
return 1;
}

void onDisconnectFailure(void* context, MQTTAsync_failureData* response)
{
printf("Disconnect failed, rc %d\\n", response->code);
disc_finished = 1;
}

void onDisconnect(void* context, MQTTAsync_successData* response)
{
printf("Successful disconnection\\n");
disc_finished = 1;
}

void onSubscribe(void* context, MQTTAsync_successData* response)
{
printf("Subscribe succeeded\\n");
subscribed = 1;
}

void onSubscribeFailure(void* context, MQTTAsync_failureData* response)
{
printf("Subscribe failed, rc %d\\n", response->code);
finished = 1;
}


void onConnectFailure(void* context, MQTTAsync_failureData* response)
{
printf("Connect failed, rc %d\\n", response->code);
finished = 1;
}


void onConnect(void* context, MQTTAsync_successData* response)
{
MQTTAsync client = (MQTTAsync)context;
MQTTAsync_responseOptions opts = MQTTAsync_responseOptions_initializer;
int rc;

printf("Successful connection\\n");

printf("Subscribing to topic %s\\nfor client %s using QoS%d\\n\\n"
"Press Q<Enter> to quit\\n\\n", TOPIC, CLIENTID, QOS);
opts.onSuccess = onSubscribe;
opts.onFailure = onSubscribeFailure;
opts.context = client;
if ((rc = MQTTAsync_subscribe(client, TOPIC, QOS, &opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start subscribe, return code %d\\n", rc);
finished = 1;
}
}


int main(int argc, char* argv[])
{
MQTTAsync client;
MQTTAsync_connectOptions conn_opts = MQTTAsync_connectOptions_initializer;
MQTTAsync_disconnectOptions disc_opts = MQTTAsync_disconnectOptions_initializer;
g_sslopts.enableServerCertAuth = 0; //for simplicity, we don't verify the server certificate
int rc;
int ch;

const char* uri = (argc > 1) ? argv[1] : ADDRESS;
g_username = (argc > 2) ? argv[2] : NULL;
g_password = (argc > 3) ? argv[3] : NULL;
printf("Using server at %s\\n", uri);

if ((rc = MQTTAsync_create(&client, uri, CLIENTID, MQTTCLIENT_PERSISTENCE_NONE, NULL))
!= MQTTASYNC_SUCCESS)
{
printf("Failed to create client, return code %d\\n", rc);
rc = EXIT_FAILURE;
goto exit;
}

if ((rc = MQTTAsync_setCallbacks(client, client, connlost, msgarrvd, NULL)) != MQTTASYNC_SUCCESS)
{
printf("Failed to set callbacks, return code %d\\n", rc);
rc = EXIT_FAILURE;
goto destroy_exit;
}

conn_opts.keepAliveInterval = 20;
conn_opts.cleansession = 1;
conn_opts.username = g_username;
conn_opts.password = g_password;
conn_opts.onSuccess = onConnect;
conn_opts.onFailure = onConnectFailure;
conn_opts.context = client;
conn_opts.ssl = &g_sslopts;
if ((rc = MQTTAsync_connect(client, &conn_opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start connect, return code %d\\n", rc);
rc = EXIT_FAILURE;
goto destroy_exit;
}

while (!subscribed && !finished)
#if defined(_WIN32)
Sleep(100);
#else
usleep(10000L);
#endif

if (finished)
goto exit;

do
{
ch = getchar();
} while (ch!='Q' && ch != 'q');

disc_opts.onSuccess = onDisconnect;
disc_opts.onFailure = onDisconnectFailure;
if ((rc = MQTTAsync_disconnect(client, &disc_opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start disconnect, return code %d\\n", rc);
rc = EXIT_FAILURE;
goto destroy_exit;
}
while (!disc_finished)
{
#if defined(_WIN32)
Sleep(100);
#else
usleep(10000L);
#endif
}

destroy_exit:
MQTTAsync_destroy(&client);
exit:
return rc;
}


/*******************************************************************************
* Copyright (c) 2023-2025 EMQ Technologies Co., William Yang and others.
*
* All rights reserved. This program and the accompanying materials
* are made available under the terms of the Eclipse Public License v2.0
* and Eclipse Distribution License v1.0 which accompany this distribution.
*
* The Eclipse Public License is available at
* https://www.eclipse.org/legal/epl-2.0/
* and the Eclipse Distribution License is available at
* http://www.eclipse.org/org/documents/edl-v10.php.
*
* MQTT over QUIC with fallback to TLS, using the serverURIs option.
*
* The client tries quic:// first and falls back to ssl:// when the QUIC
* connection cannot be established (e.g. UDP is blocked). Note that a
* blocked/dead QUIC port fails by timeout, not immediately: expect up to
* ~30 seconds per URI before the fallback kicks in (bounded by
* connectTimeout and the QUIC handshake timeout), and with
* MQTTVERSION_DEFAULT each URI is retried once per MQTT protocol version,
* so setting an explicit version halves the failover time.
*/

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include "MQTTAsync.h"

#if !defined(_WIN32)
#include <unistd.h>
#else
#include <windows.h>
#endif

#if defined(_WRS_KERNEL)
#include <OsWrapper.h>
#endif

#define QUIC_ADDRESS "quic://127.0.0.1:14567"
#define SSL_ADDRESS "ssl://127.0.0.1:8883"
#define CLIENTID "ExampleClientPubFallback"
#define TOPIC "MQTT Examples"
#define PAYLOAD "Hello World!"
#define QOS 1
#define TIMEOUT 10000L

int finished = 0;

void connlost(void *context, char *cause)
{
printf("\\nConnection lost\\n");
printf(" cause: %s\\n", cause);
finished = 1;
}

void onDisconnectFailure(void* context, MQTTAsync_failureData* response)
{
printf("Disconnect failed\\n");
finished = 1;
}

void onDisconnect(void* context, MQTTAsync_successData* response)
{
printf("Successful disconnection\\n");
finished = 1;
}

void onSendFailure(void* context, MQTTAsync_failureData* response)
{
printf("Message send failed token %d error code %d\\n", response->token, response->code);
finished = 1;
}

void onSend(void* context, MQTTAsync_successData* response)
{
MQTTAsync client = (MQTTAsync)context;
MQTTAsync_disconnectOptions opts = MQTTAsync_disconnectOptions_initializer;
int rc;

printf("Message with token value %d delivery confirmed\\n", response->token);
opts.onSuccess = onDisconnect;
opts.onFailure = onDisconnectFailure;
opts.context = client;
if ((rc = MQTTAsync_disconnect(client, &opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start disconnect, return code %d\\n", rc);
exit(EXIT_FAILURE);
}
}

void onConnectFailure(void* context, MQTTAsync_failureData* response)
{
printf("Connect failed on all serverURIs, rc %d\\n", response ? response->code : 0);
finished = 1;
}

void onConnect(void* context, MQTTAsync_successData* response)
{
MQTTAsync client = (MQTTAsync)context;
MQTTAsync_responseOptions opts = MQTTAsync_responseOptions_initializer;
MQTTAsync_message pubmsg = MQTTAsync_message_initializer;
int rc;

printf("Successful connection via %s\\n", response->alt.connect.serverURI);
opts.onSuccess = onSend;
opts.onFailure = onSendFailure;
opts.context = client;
pubmsg.payload = PAYLOAD;
pubmsg.payloadlen = (int)strlen(PAYLOAD);
pubmsg.qos = QOS;
pubmsg.retained = 0;
if ((rc = MQTTAsync_sendMessage(client, TOPIC, &pubmsg, &opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start sendMessage, return code %d\\n", rc);
exit(EXIT_FAILURE);
}
}

int messageArrived(void* context, char* topicName, int topicLen, MQTTAsync_message* m)
{
/* not expecting any messages */
return 1;
}

int main(int argc, char* argv[])
{
MQTTAsync client;
MQTTAsync_connectOptions conn_opts = MQTTAsync_connectOptions_initializer;
MQTTAsync_SSLOptions sslopts = MQTTAsync_SSLOptions_initializer;
sslopts.enableServerCertAuth = 0; //for simplicity, we don't verify the server certificate
char* uris[2];
int rc;

const char* quic_uri = (argc > 1) ? argv[1] : QUIC_ADDRESS;
const char* ssl_uri = (argc > 2) ? argv[2] : SSL_ADDRESS;
const char* username = (argc > 3) ? argv[3] : NULL;
const char* password = (argc > 4) ? argv[4] : NULL;

uris[0] = (char*)quic_uri; /* tried first */
uris[1] = (char*)ssl_uri; /* fallback */

if ((rc = MQTTAsync_create(&client, quic_uri, CLIENTID, MQTTCLIENT_PERSISTENCE_NONE, NULL)) != MQTTASYNC_SUCCESS)
{
printf("Failed to create client object, return code %d\\n", rc);
exit(EXIT_FAILURE);
}

if ((rc = MQTTAsync_setCallbacks(client, NULL, connlost, messageArrived, NULL)) != MQTTASYNC_SUCCESS)
{
printf("Failed to set callback, return code %d\\n", rc);
exit(EXIT_FAILURE);
}

conn_opts.keepAliveInterval = 20;
conn_opts.cleansession = 1;
conn_opts.username = username;
conn_opts.password = password;
conn_opts.serverURIs = uris;
conn_opts.serverURIcount = 2;
conn_opts.MQTTVersion = MQTTVERSION_3_1_1; /* avoid a second attempt per URI with MQTT 3.1 */
conn_opts.connectTimeout = 10; /* each URI is tried for up to 10s */
conn_opts.onSuccess = onConnect;
conn_opts.onFailure = onConnectFailure;
conn_opts.context = client;
conn_opts.ssl = &sslopts;
if ((rc = MQTTAsync_connect(client, &conn_opts)) != MQTTASYNC_SUCCESS)
{
printf("Failed to start connect, return code %d\\n", rc);
exit(EXIT_FAILURE);
}

printf("Trying %s, falling back to %s\\n", quic_uri, ssl_uri);
while (!finished)
#if defined(_WIN32)
Sleep(100);
#else
usleep(10000L);
#endif

MQTTAsync_destroy(&client);
return rc;
}



Java 示例代码

QuicExample
QuicPublish
QuicSubscribe
/*
* Copyright 2018-present HiveMQ and the HiveMQ Community
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package com.hivemq.client.mqtt.examples;

import com.hivemq.client.mqtt.mqtt5.Mqtt5BlockingClient;
import com.hivemq.client.mqtt.mqtt5.Mqtt5Client;
import com.hivemq.client.mqtt.mqtt5.Mqtt5ClientBuilder;

import java.nio.charset.StandardCharsets;

/**
* Shared setup for the QUIC publish/subscribe examples.
*/
final class QuicExample {

static Mqtt5BlockingClient connect(final String identifier, final String[] args) {
final String server = (args.length > 0) ? args[0] : "localhost";
final Mqtt5ClientBuilder builder = Mqtt5Client.builder().identifier(identifier).quicWithDefaultConfig();
applyServer(builder, server);

final Mqtt5BlockingClient client = builder.buildBlocking();
if (args.length >= 3) {
client.connectWith()
.simpleAuth()
.username(args[1])
.password(args[2].getBytes(StandardCharsets.UTF_8))
.applySimpleAuth()
.send();
} else {
client.connect();
}
System.out.println("connected over QUIC to " + server);
return client;
}

private static void applyServer(final Mqtt5ClientBuilder builder, final String server) {
final int colon = server.lastIndexOf(':');
if ((colon > 0) && (server.indexOf(']') < colon)) {
builder.serverHost(server.substring(0, colon)).serverPort(Integer.parseInt(server.substring(colon + 1)));
} else {
builder.serverHost(server);
}
}

private QuicExample() {}
}

/*
* Copyright 2018-present HiveMQ and the HiveMQ Community
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package com.hivemq.client.mqtt.examples;

import com.hivemq.client.mqtt.datatypes.MqttQos;
import com.hivemq.client.mqtt.mqtt5.Mqtt5BlockingClient;

import java.nio.charset.StandardCharsets;

/**
* Publishes to {@code demo/quic} over MQTT over QUIC.
* <p>
* Usage: {@code QuicPublish [host[:port]] [username] [password]}
* <p>
* Run {@link QuicSubscribe} in another process, then start this example. Requires a broker that speaks MQTT over
* QUIC (TLS 1.3, ALPN {@code mqtt}). If no host is given, {@code localhost:14567} is used. Add the
* {@code hivemq-mqtt-client-quic} module to the classpath.
*
* @author HiveMQ
*/
public class QuicPublish {

public static void main(final String[] args) {
final Mqtt5BlockingClient client = QuicExample.connect("quic-publish-example", args);
try {
for (int i = 0; i < 5; i++) {
final String payload = "hello quic " + i;
client.publishWith()
.topic("demo/quic")
.qos(MqttQos.AT_LEAST_ONCE)
.payload(payload.getBytes(StandardCharsets.UTF_8))
.send();
System.out.println("published over QUIC: " + payload);
}
} finally {
client.disconnect();
}
}
}

/*
* Copyright 2018-present HiveMQ and the HiveMQ Community
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package com.hivemq.client.mqtt.examples;

import com.hivemq.client.mqtt.MqttGlobalPublishFilter;
import com.hivemq.client.mqtt.datatypes.MqttQos;
import com.hivemq.client.mqtt.mqtt5.Mqtt5BlockingClient;
import com.hivemq.client.mqtt.mqtt5.Mqtt5BlockingClient.Mqtt5Publishes;

/**
* Subscribes to {@code demo/quic} over MQTT over QUIC.
* <p>
* Usage: {@code QuicSubscribe [host[:port]] [username] [password]}
* <p>
* Requires a broker that speaks MQTT over QUIC (TLS 1.3, ALPN {@code mqtt}). If no host is given,
* {@code localhost:14567} is used. Add the {@code hivemq-mqtt-client-quic} module to the classpath.
*
* @author HiveMQ
*/
public class QuicSubscribe {

public static void main(final String[] args) throws InterruptedException {
final Mqtt5BlockingClient client = QuicExample.connect("quic-subscribe-example", args);
try (final Mqtt5Publishes publishes = client.publishes(MqttGlobalPublishFilter.ALL)) {
client.subscribeWith().topicFilter("demo/quic").qos(MqttQos.AT_LEAST_ONCE).send();
System.out.println("subscribed over QUIC, waiting for messages on demo/quic ...");
while (!Thread.currentThread().isInterrupted()) {
System.out.println(publishes.receive());
}
} finally {
client.disconnect();
}
}
}