跳转到内容
新建笔记

Qt TCP:字节流分帧、回环验证与线程归属

Qt 使用 QTcpSocket 表示一条 TCP 连接,使用 QTcpServer 监听并接受连接。客户端主动连接,服务端监听和接受;客户端流程里不包含 listen() 或 accept()。

TCP 提供有序字节流,不保留应用消息的边界。一次 readyRead 可能只带来半条消息,也可能包含多条消息;readAll() 只是取走当前可读字节。需要先定义应用层分帧,再将完整消息交给界面或业务代码。

连接与信号的对应关系

跳转到“连接与信号的对应关系”
阶段客户端服务端
创建创建有明确所有者的 QTcpSocket创建 QTcpServer
建立先连接信号,再调用 connectToHost(host, port)检查 listen(address, port) 的结果
就绪connected 表示连接成功newConnection 后循环取 nextPendingConnection()
收数据在 readyRead 中读取并交给分帧器对每条已接受连接分别处理 readyRead
错误与结束errorOccurred 提供错误;disconnected 表示断开监听错误与已连接 socket 的错误分别处理

原片段的 new QTcpSocket 不能赋给 QTcpServer*;nextPendingConnection、readyRead、QHostAddress 的拼写也必须准确。消息框头文件是 QMessageBox,不是 QMessage。disconnected 不等于“从未连接成功”,例如对端正常关闭也会触发它。QTcpServer、QAbstractSocket

端口类型是 quint16。从文本取得端口时,先用带成功标志的整数转换,检查业务范围再转换类型;不能用有符号 short 或直接强转吞掉溢出。示例让服务端监听本机回环地址、使用端口 0 由系统分配空闲端口,再通过 serverPort() 获取实际端口。

下面定义一个小型行协议:以 LF 字节 \n 结束一帧;一帧最多 1024 字节,内容不能包含 LF,空帧合法。内容按字节积累,完整后才可按约定的 UTF-8 等编码解释。二进制负载需要另选长度前缀等协议,不能直接沿用这个限制。

示例的服务端在主线程处理回显,客户端 socket 在工作线程创建和使用。这样同时展示原笔记的客户端、服务端和多线程用途。普通低负载 TCP 通信本就可以使用事件驱动方式,不需要每条连接都建立线程。

完整程序:拆帧、合帧与实际回环通信

跳转到“完整程序:拆帧、合帧与实际回环通信”

使用 Qt 6、C++17,链接 Qt6::Core 与 Qt6::Network。程序自动结束,Debug 配置中的断言验证分帧边界;不连接外部服务。

#include <QAbstractSocket>
#include <QByteArray>
#include <QCoreApplication>
#include <QDebug>
#include <QHostAddress>
#include <QList>
#include <QMetaObject>
#include <QObject>
#include <QPointer>
#include <QString>
#include <QTcpServer>
#include <QTcpSocket>
#include <QThread>
#include <QTimer>
#include <cassert>
#include <memory>
class LineDecoder {
public:
static constexpr qsizetype MaximumFrame = 1024;
bool feed(const QByteArray& chunk, QList<QByteArray>& frames) {
if (failed_ || chunk.size() > 4096) {
failed_ = true;
return false;
}
pending_ += chunk;
for (;;) {
const qsizetype end = pending_.indexOf('\n');
if (end < 0)
break;
if (end > MaximumFrame) {
failed_ = true;
return false;
}
frames.append(pending_.left(end));
pending_.remove(0, end + 1);
}
if (pending_.size() > MaximumFrame) {
failed_ = true;
return false;
}
return true;
}
bool atBoundary() const { return !failed_ && pending_.isEmpty(); }
private:
QByteArray pending_;
bool failed_ = false;
};
bool enqueue(QTcpSocket* socket, const QByteArray& bytes) {
if (socket->state() != QAbstractSocket::ConnectedState ||
socket->bytesToWrite() + bytes.size() > 64 * 1024)
return false;
// 本例将部分接收或错误视为失败并终止,绝不忽略未接受的后缀。
return socket->write(bytes) == bytes.size();
}
void checkFraming() {
LineDecoder decoder;
QList<QByteArray> frames;
assert(decoder.feed("he", frames) && frames.isEmpty());
assert(!decoder.atBoundary());
assert(decoder.feed("llo\nworld\n\n", frames));
assert((frames == QList<QByteArray>{"hello", "world", ""}));
assert(decoder.atBoundary());
LineDecoder tooLong;
frames.clear();
assert(!tooLong.feed(QByteArray(1025, 'x'), frames));
LineDecoder limit;
frames.clear();
assert(limit.feed(QByteArray(1024, 'x') + '\n', frames));
assert(frames.size() == 1 && frames.first().size() == 1024);
const QByteArray utf8 = QStringLiteral("你好").toUtf8();
LineDecoder unicode;
frames.clear();
assert(unicode.feed(utf8.left(1), frames) && frames.isEmpty());
assert(unicode.feed(utf8.mid(1) + '\n', frames));
assert(QString::fromUtf8(frames.first()) == QStringLiteral("你好"));
}
int main(int argc, char* argv[]) {
QCoreApplication app(argc, argv);
checkFraming();
bool completed = false;
bool success = false;
const auto finish = [&](bool ok, const QString& message) {
// 可能从工作线程调用;结果状态和退出动作统一回到主线程。
QMetaObject::invokeMethod(&app, [&, ok, message] {
if (completed)
return;
completed = true;
success = ok;
if (!ok)
qWarning().noquote() << message;
app.quit();
}, Qt::QueuedConnection);
};
QTcpServer server;
QPointer<QTcpSocket> peer;
QObject::connect(&server, &QTcpServer::newConnection, &app, [&] {
while (server.hasPendingConnections()) {
QTcpSocket* socket = server.nextPendingConnection();
if (peer) { // 本演示只接受自己的一个客户端
socket->abort();
socket->deleteLater();
continue;
}
peer = socket;
socket->setReadBufferSize(64 * 1024);
auto decoder = std::make_shared<LineDecoder>();
QObject::connect(socket, &QTcpSocket::readyRead, socket,
[socket, decoder, finish] {
while (socket->bytesAvailable() > 0) {
QList<QByteArray> frames;
if (!decoder->feed(socket->read(4096), frames)) {
socket->abort();
finish(false, "Server rejected an oversized frame");
return;
}
for (const QByteArray& frame : frames) {
if (!enqueue(socket, frame + '\n')) {
socket->abort();
finish(false, "Server write was not fully accepted");
return;
}
}
}
});
QObject::connect(socket, &QTcpSocket::disconnected, socket,
[socket, decoder, finish] {
if (!decoder->atBoundary())
finish(false, "Connection ended inside a frame");
socket->deleteLater();
});
QObject::connect(socket, &QTcpSocket::errorOccurred, socket,
[socket, finish](QAbstractSocket::SocketError error) {
if (error != QAbstractSocket::RemoteHostClosedError)
finish(false, socket->errorString());
});
}
});
if (!server.listen(QHostAddress::LocalHost, 0)) {
qWarning() << server.errorString();
return 1;
}
const quint16 port = server.serverPort();
QThread thread;
auto* worker = new QObject;
worker->moveToThread(&thread);
QObject::connect(&thread, &QThread::finished, worker, &QObject::deleteLater);
QObject::connect(&thread, &QThread::started, worker, [worker, port, finish] {
auto* socket = new QTcpSocket(worker); // 此行在工作线程执行
assert(socket->thread() == QThread::currentThread());
socket->setReadBufferSize(64 * 1024);
auto decoder = std::make_shared<LineDecoder>();
auto received = std::make_shared<QList<QByteArray>>();
QObject::connect(socket, &QTcpSocket::connected, socket, [socket, finish] {
if (!enqueue(socket, "he")) {
finish(false, "Client write failed");
return;
}
QTimer::singleShot(20, socket, [socket, finish] {
if (!enqueue(socket, "llo\nworld\n"))
finish(false, "Client write failed");
});
});
QObject::connect(socket, &QTcpSocket::readyRead, socket,
[socket, decoder, received, finish] {
while (socket->bytesAvailable() > 0) {
QList<QByteArray> frames;
if (!decoder->feed(socket->read(4096), frames)) {
finish(false, "Client rejected an oversized frame");
return;
}
received->append(frames);
if (received->size() >= 2) {
const bool correct = *received == QList<QByteArray>{"hello", "world"};
finish(correct && decoder->atBoundary(), "Unexpected echo frames");
return;
}
}
});
QObject::connect(socket, &QTcpSocket::errorOccurred, socket,
[socket, finish](QAbstractSocket::SocketError) {
finish(false, socket->errorString());
});
QObject::connect(socket, &QTcpSocket::disconnected, socket, [finish] {
finish(false, "Client disconnected before completion");
});
socket->connectToHost(QHostAddress::LocalHost, port);
});
QTimer::singleShot(3000, &app, [finish] { finish(false, "Loopback check timed out"); });
thread.start();
app.exec();
thread.quit();
thread.wait(); // worker 与其 socket 在工作线程结束时清理
if (peer) {
peer->abort();
delete peer.data();
}
server.close();
return success ? 0 : 1;
}

分帧器的直接断言确定检查了拆帧、连续多帧、空帧、长度上限和 UTF-8 字节拆分。实际网络部分验证连接、两条回显和线程退出;操作系统仍可能合并两次发送,所以不能仅凭两次 write() 就声称网络必然分成两次接收。

write() 成功表示字节被写入流程接受,不表示对端业务已经处理完成。本例以收到两条完整回显作为应用层成功条件。持续传输的产品还应设计待发队列、流量控制、超时和每次事件处理量;本例超过发送队列上限或出现部分接受就明确终止,没有实现重试协议。

为什么原来的 QThread 片段不能工作

跳转到“为什么原来的 QThread 片段不能工作”

把主线程的 QTcpSocket* 传进一个 QThread 子类构造函数,不会改变 socket 的线程归属。QThread 对象本身通常仍属于创建它的线程,其成员槽也不会因为写在该类里就自动跑到 run() 的线程。

原 run() 只建立连接便返回,没有进入事件循环。事件驱动的 socket 无法以预期方式持续处理通知。上例将 worker 移入线程,再在 worker 的线程内创建 socket,并使用 QThread 默认事件循环;所有 socket 回调都有该 socket 作为上下文。线程与 QObject

另一种服务端方案是重写 QTcpServer::incomingConnection(qintptr),将原生描述符交给目标线程,在目标线程创建 QTcpSocket 并调用 setSocketDescriptor()。不能让两个 socket 同时拥有同一描述符,也要处理采用失败时的关闭责任。该方案需要独立实现,不能直接把已经有父对象的已接受 socket 指针当作可任意跨线程访问的资源。incomingConnection

若收到数据后更新 QLineEdit 等 GUI 控件,应向 GUI 线程发送排队信号或调用,保持控件在 GUI 线程访问。线程中的网络处理也不应依赖对 sender() 的未检查 C 风格强转;明确捕获连接对象更容易看清归属与寿命。