惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Martin Fowler
Martin Fowler
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
IT之家
IT之家
美团技术团队
酷 壳 – CoolShell
酷 壳 – CoolShell
Y
Y Combinator Blog
T
Tailwind CSS Blog
D
Docker
博客园 - Franky
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
Google DeepMind News
Google DeepMind News
腾讯CDC
Vercel News
Vercel News
Engineering at Meta
Engineering at Meta
U
Unit 42
The Cloudflare Blog
S
SegmentFault 最新的问题
WordPress大学
WordPress大学
爱范儿
爱范儿
Recent Announcements
Recent Announcements
博客园 - 聂微东
博客园 - 叶小钗
H
Help Net Security
MyScale Blog
MyScale Blog

博客园 - delphi中间件

delphi 面向模型编程 array of TVarRec core.recordModel.pas 使用泛型序列结构体 DDD建模指导 rabbitMQ VS mqtt redis流的应用场景 redis流的操作命令 领域服务与领域事件 业务规则和模型 限界上下文与统一语言 领域驱动 mqtt即时通讯 ActiveRecord ORM RAD(速成应用开发) unigui插件框架 工厂流水线式自动生产UNIGUI WEB软件 delphi cs\web一种统一的界面风格 - delphi中间件 - 博客园 动态生成unidbgrid 单据工厂 用json元数据填充模板 mormot2 ORM rest vs jsonrpc SSE技术详解:使用 HTTP 做服务端数据推送应用的技术 http持久连接 json-rpc 2.0 MCP服务器 RTTI对性能的影响 频繁地创建和销毁对象 TMultiPartFormData
redis消费者组
delphi中间件 · 2026-07-28 · via 博客园 - delphi中间件

消费组特性:
消费组是一组消费者,它们可以协同消费流中的事件。每个消费者组都有一个消费者组名称,用于标识不同的消费者组。
消费组中的每个消费者都有一个消费者名称,用于标识不同的消费者。
消费组会维护每个消费者的消费状态,以确保每个事件只被一个消费者处理
Redis流还支持阻塞消费,即消费者可以使用XREADGROUP命令来等待新的事件,并在事件到达时立即处理它们。
流支持对事件设置字段和值,支持时间戳,允许灵活的数据建模。

使用delphiredisclient开源控件。

uses Redis.Commons, Redis.Client, Redis.Values, Redis.NetLib.Indy;

//XGROUP CREATE mystream mygroup 0,从流头开始
//只关心实时消息:用 XGROUP CREATE mystream mygroup $(推荐默认做法)
//已有 Group 想让新消费者补读:必须用 XREADGROUP GROUP mygroup newconsumer COUNT 100 STREAMS mystream 0,显式指定起始 ID 0;但注意这会干扰 Group 整体的 pending 状态管理

  //创建消费者组
  var lCmd3: IRedisCommand := NewRedisCommand('XGROUP');
  lCmd3
    .Add('create')
    .Add('mystream')   //流名
    .Add('mygroup')    //组名
    .Add('$')
    .Add('MKSTREAM');
  var lRes3: TRedisNullable<string> := fRedis.ExecuteWithStringResult(lCmd3);

  //为消费者组创建消费者
  var lCmd1: IRedisCommand := NewRedisCommand('XGROUP');
  lCmd1
    .Add('CREATECONSUMER')
    .Add('mystream')    //流名
    .Add('mygroup')     //组名
    .Add('myconsumer'); //消费者名
  var lRes1: TRedisNullable<string> := fRedis.ExecuteWithStringResult(lCmd1);

  //往流里插入一条消息
  var lCmd: IRedisCommand := NewRedisCommand('XADD');
  lCmd
    .Add('mystream')
    .Add('MAXLEN')
    .Add('~')
    .Add(10)
    .Add('*')
    .Add('key1')     //键1
    .Add('value1');  //值1
  var lRes: TRedisNullable<string> := fRedis.ExecuteWithStringResult(lCmd);
  Log(lRes.Value);
procedure TMainForm.Button2Click(Sender: TObject);
//订阅消费者组
begin
  TTask.Run(
    procedure
    begin
      while TTask.CurrentTask.Status <> TTaskStatus.Canceled do
      begin
        var lRedis := NewRedisClient();
        var lCmd: IRedisCommand := NewRedisCommand('XREADGROUP');
        lCmd
          .Add('group')
          .Add('mygroup')
          .Add('myconsumer')
          .Add('block')
          .Add(5000)
          .Add('count')
          .Add(1)
          .Add('streams')
          .Add('mystream')
          .Add('>');
        var lRes: TRedisRESPArray := lRedis.ExecuteAndGetRESPArray(lCmd);
        if Assigned(lRes) then
        Log(lRes.tostring);
      end;
    end);
end;

获取REDIS命令执行结果

        var lRes: TRedisRESPArray := lRedis.ExecuteAndGetRESPArray(lCmd);
        if not Assigned(lRes) then
        begin
          Continue;
        end;
        try
          if lRes.Count > 0 then
          begin
           // var
            lSizeOfMyStreamArray := lRes
              .Items[0].ArrayValue
              .Items[1].ArrayValue
              .Count;

            lLastID := lRes
              .Items[0].ArrayValue
              .Items[1].ArrayValue
              .Items[lSizeOfMyStreamArray - 1].ArrayValue
              .Items[0].Value;
          end;
          Log(lRes.ToString+' count:'+lSizeOfMyStreamArray.ToString+' lastid:'+lLastID);
        finally
          lRes.Free;
        end;

XPENDING 命令用于查询流中的延时消息。它可以返回关于流中哪些消息是未处理的(即尚未被消费者处理的),以及这些消息的一些统计信息,如消息的数量、最小的 ID 和最大的 ID 等。

XACK‌用于确认 Redis Stream 消费者组中的消息已处理成功,将其从待处理列表(PEL)移除

uses
  redis.Command, redis.Client;

var
  Redis: IRedisClient;
  Cmd: IRedisCommand;
  AckCount: Integer;
begin
  Redis := NewRedisClient('127.0.0.1', 6379);
  
  // 方式一:直接构建命令(推荐,兼容性好)
  Cmd := NewRedisCommand('XACK')
    .Add('mystream')       // Stream Key
    .Add('mygroup')        // Consumer Group
    .Add('1526569495631-0'); // 消息 ID(可追加多个)
  
  AckCount := Redis.ExecuteWithIntegerResult(Cmd); 
  // 结果:成功确认的消息数
  
  // 方式二:批量确认多个 ID
  Cmd := NewRedisCommand('XACK')
    .Add('mystream')
    .Add('mygroup')
    .Add('1526569495631-0')
    .Add('1526569495632-0');
  AckCount := Redis.ExecuteWithIntegerResult(Cmd);
end