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

使用IdTCPServer和TIdThreadSafeObjectList发送TMemoryStream遇空流问题

基于TMemoryStream的Indy TCP服务器数据发送问题

调用TMyContext.SendQueue()向客户端发送数据时流为空,同时怀疑SendQueue()未将流发送到正确客户端,是否需要遍历所有上下文?


简化后的问题代码

TMyContext = class(TIdServerContext)
  private
    Context_ID: string;
    Queue: TIdThreadSafeObjectList;
    QueuePending: Boolean;
  public
    constructor Create(AConnection: TIdTCPConnection; AYarn: TIdYarn; AList: TIdContextThreadList = nil); override;
    destructor Destroy; override;
    procedure AddToQueue(const ms: TMemoryStream);
    procedure SendQueue;
  end;
...
...

constructor TMyContext.Create(AConnection: TIdTCPConnection; AYarn: TIdYarn; AList: TIdContextThreadList = nil);
begin
  inherited;
  Queue := TIdThreadSafeObjectList.Create;
end;

destructor TMyContext.Destroy;
begin
  Queue.Free;
  inherited;
end;

procedure TMyContext.AddToQueue(const ms: TMemoryStream);
var
  List: TList;
begin
  List := Queue.LockList;
  try
    list.Add(ms);
    QueuePending := True;
  finally
    Queue.UnlockList;
  end;
end;

procedure TMyContext.SendQueue;
var
  list: TList;
  i: Integer;
begin
  if not QueuePending then
    Exit;

  list := Queue.LockList;
  try
    if list.Count = 0 then
    begin
      QueuePending := False;
      Exit;
    end;

    for i := 0 to List.Count-1 do
      Connection.IOHandler.Write( TMemoryStream(List[i]), 0, True); // 此处List[i]大小为0
      // OnExecute()中的TMemoryStream对象是否仍存在?

    list.Clear;
    QueuePending := False;
  finally
    Queue.UnlockList;
  end;
end;
...
...
procedure TServer.IdTCPServerConnect(AContext: TIdContext);
var
  LContext: TMyContext;
  ms: TMemorySTream;
begin
  LContext := TMyContext(AContext);

  ms := TMemoryStream.Create;
  try
    ...
    ...
    AContext.Connection.IOHandler.ReadStream(msgFromClient, size);
    msgFromClient.Position:= 0;
    ...

    // 从流中获取一些值
    ...
    LContext.Context_ID  := Get_ID(ms);

  finally
    ms.Free;
  end;
end;

procedure TServer.IdTCPServerExecute(AContext: TIdContext);
var
  ms: TMemoryStream;
  LContext,Ctx_Peer2: TMyContext;
  List: TList;
  I: integer;
begin
  LContext := TMyContext(AContext);

  // 发送服务器中待处理的消息
  LContext.SendQueue;


  // 检查是否收到客户端消息
  if AContext.Connection.IOHandler.InputBufferIsEmpty then
  begin
    AContext.Connection.IOHandler.CheckForDataOnSource(100);
    AContext.Connection.IOHandler.CheckForDisconnect;
    if AContext.Connection.IOHandler.InputBufferIsEmpty then Exit;
  end;

  size:=AContext.Connection.IOHandler.ReadLongInt;
  try
    try
      ms := TMemoryStream.Create;
      AContext.Connection.IOHandler.ReadStream(msgStream, size);
    except
      Exit;
    end;

    // 对ms流执行一些操作
    ...
    ...

    // 将新数据写入ms并添加至调用上下文队列
    ms.clear:= 0;
    ...
    ...
    ms.position := 0;
    LContext.AddToQueue(ms);

    // 将新数据写入ms并发送至另一对等端
    // 此处是否需要使用与ms不同的TMemoryStream变量?
    ms.clear:= 0;
    ...
    ...
    ms.position := 0;
    List := IdTCPServer.Contexts.LockList;
    try

      for I := 0 to List.Count-1 do
      begin
        if (TMyContext(List[I].Context_ID = context_id_i_want_to_send) then
        begin
          Ctx_Peer2 := TMyContext(List[I]);
          Break;
        end;
      end;
      Ctx_Peer2.AddToQueue(ms);
    finally
      IdTCPServer.Contexts.UnlockList;
    end;

  finally
    ms.Free;
  end;
end;

procedure TServer.IdTCPServerDisconnect(AContext: TIdContext);
var
  Ctx: TMyContext;
begin
  Ctx := TMyContext(AContext);
  Ctx.Context_ID = '';
  Ctx.Queue.Clear;
end;

核心问题分析

1. 流为空的直接原因

  • 重复复用+提前释放流:IdTCPServerExecute中复用同一个ms添加到两个队列,最后在finally里直接Free,导致SendQueue执行时队列中的流已经被释放,或因重复Clear丢失数据。
  • 对象所有权冲突:TIdThreadSafeObjectList默认接管对象所有权(FreeOnRelease=True),但你在外部又手动释放了加入队列的流,造成野指针或双重释放,读取时流为空。

2. SendQueue的发送目标问题

SendQueue是每个TMyContext独立的方法,只会向当前Context对应的客户端发送队列数据,逻辑是正确的——不需要遍历所有上下文,每个客户端的待发队列本就应该独立处理。


修复方案

针对流为空的问题

  1. 为每个队列项创建独立流副本
    向多个队列添加数据时,必须创建独立的流副本,不能复用同一个流:

    // 添加到自身队列
    var msCopy1 := TMemoryStream.Create;
    msCopy1.CopyFrom(ms, 0);
    msCopy1.Position := 0;
    LContext.AddToQueue(msCopy1);
    
    // 添加到对等端队列
    var msCopy2 := TMemoryStream.Create;
    msCopy2.CopyFrom(ms, 0);
    msCopy2.Position := 0;
    Ctx_Peer2.AddToQueue(msCopy2);
    
  2. 调整列表的对象所有权设置

    • 若不想让列表自动释放对象,创建时设置FreeOnRelease=False,并在发送后手动释放流:
      // 创建列表时关闭自动释放
      Queue := TIdThreadSafeObjectList.Create(False);
      
      // SendQueue中发送后释放流
      for i := 0 to List.Count-1 do
      begin
        var stream := TMemoryStream(List[i]);
        Connection.IOHandler.Write(stream, 0, True);
        stream.Free;
      end;
      list.Clear;
      
    • 若保持默认FreeOnRelease=True,则不要在Execute的finally里释放加入队列的流,只释放未加入队列的临时流。

针对SendQueue的逻辑优化

  • QueuePending必须在锁内修改,避免多线程竞争:
    list := Queue.LockList;
    try
      if list.Count = 0 then
      begin
        QueuePending := False;
        Exit;
      end;
      // ...发送逻辑
      list.Clear;
      QueuePending := False;
    finally
      Queue.UnlockList;
    end;
    
  • 确保发送前流的Position为0,避免从流的末尾开始发送空数据。

其他代码问题修复

  • IdTCPServerDisconnect中赋值用:=而非=:Ctx.Context_ID := '';
  • IdTCPServerExecute中读取流的变量名错误:ReadStream(msgStream, size)应改为ReadStream(ms, size)
  • 遍历Contexts时要处理Ctx_Peer2未找到的情况,避免空指针:
    Ctx_Peer2 := nil;
    for I := 0 to List.Count-1 do
    begin
      var ctx := TMyContext(List[I]);
      if ctx.Context_ID = context_id_i_want_to_send then
      begin
        Ctx_Peer2 := ctx;
        Break;
      end;
    end;
    if Assigned(Ctx_Peer2) then
      Ctx_Peer2.AddToQueue(msCopy2);
    

内容的提问来源于stack exchange,提问作者R.Schirru

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 14:20:18