-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathQueue.php
More file actions
1305 lines (1187 loc) · 55.5 KB
/
Copy pathQueue.php
File metadata and controls
1305 lines (1187 loc) · 55.5 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
<?php
namespace TypechoPlugin\Access;
use Redis;
use Typecho\Config;
use Typecho\Db;
if (!defined('__TYPECHO_ROOT_DIR__')) {
exit;
}
/**
* 访问日志写入队列
*
* 把每次访问先塞进 Redis 列表,攒够一批再用一条多值 INSERT 落库,
* 目的有两个:
* - 削掉每次访问的建连接开销(独立数据库模式下这一项占单次写入的六成)
* - 把数据库连接数从「每次访问一条」降到「每批一条」,避免突发流量打满 max_connections
*
* 消息在 Redis 里经过三个列表:主队列 -> processing(正在写库)-> 死信(写不进去的)。
* 中间那一步不是多余的:直接「读了再按位置裁掉」的话,生产者的裁剪一旦插进来,
* 消费者就会裁掉别人刚写进来的消息。
*
* 刷库由请求顺带触发:达到条数或时间阈值时,本次请求抢到锁的那一个负责刷,
* 并且推迟到响应发出之后执行,访客感知不到。
* 另外控制台加载数据时会同步刷一次,命令行脚本可挂 cron 兜底。
*
* 没有 Redis 时整套机制不启用,写入行为与之前完全一致。
*/
final class Queue
{
/**
* 队列相关键的名字(不含前缀)
*
* 完整键名由 Cache::key() 补上带站点指纹的前缀,所以这里全部是方法而不是常量:
* 常量在编译期定型,写死的话多个站点共用一个 Redis 就会共用同一条队列。
*/
private const NAME = 'queue';
private const NAME_PROCESSING = 'queue:processing';
private const NAME_DEAD = 'queue:dead';
private const NAME_LOCK = 'queue:lock';
private const NAME_LAST_FLUSH = 'queue:last_flush';
/**
* 队首这一批「从什么时候开始就写不进去」的时间戳
*
* 只有把整批留下重试时才会写它,批次一旦确认掉就删。
* 用途见 STUCK_SECONDS。
*/
private const NAME_STUCK_SINCE = 'queue:stuck_since';
/**
* 加站点指纹之前用过的固定键名(旧前缀下的第一代)
*
* 升级之后队列会换到新键名上,这几个键里可能还压着没落库的访问日志。
* adoptLegacy() 负责把它们接管过来,不接管就等于丢数据。
*
* 注意旧前缀下还有第二代(typecho_access:{指纹}:queue*,加了指纹但没换前缀),
* 那一代按 v3.2.3 的决定**不做接管**,只在卸载清理时被保护和识别,
* 见 isLegacyDataKey() 与 adoptLegacy() 的说明。
*/
private const LEGACY_PREFIX = Cache::LEGACY_BASE . 'queue';
/** 待写入队列(Redis List) */
public static function key(): string
{
return Cache::key(self::NAME);
}
/**
* 正在写库的一批(Redis List)
*
* 消费者不再「读了之后按位置裁掉」,而是用一段 Lua 把消息原子地从主队列搬到这里,
* 写库成功再从这里清掉。中途崩掉的话数据留在这里,下一次刷库会先把它捡回来。
*/
public static function processingKey(): string
{
return Cache::key(self::NAME_PROCESSING);
}
/**
* 死信队列(Redis List)
*
* 解析不了、或者数据库明确拒绝的消息落到这里,而不是直接丢掉。
* 队列的确认是整批 LTRIM,做不到「只确认成功的那几条」(Redis List 不支持按位置挑着删),
* 所以退而求其次:裁掉之前先把失败的原样留一份证据,可以人工排查或改完再回放。
*/
public static function deadKey(): string
{
return Cache::key(self::NAME_DEAD);
}
/** 刷库互斥锁 */
private static function lockKey(): string
{
return Cache::key(self::NAME_LOCK);
}
/** 上次刷库时间 */
private static function lastFlushKey(): string
{
return Cache::key(self::NAME_LAST_FLUSH);
}
/** 队首批次卡住的起始时间 */
private static function stuckSinceKey(): string
{
return Cache::key(self::NAME_STUCK_SINCE);
}
/**
* 记下「队首这批从现在开始卡住了」,已经记过就不覆盖
*
* @param Redis $redis
* @return void
*/
private static function markStuck(Redis $redis): void
{
try {
# NX:第一次卡住的时间才是起点,后面每轮重试都覆盖的话永远到不了上限
$redis->set(self::stuckSinceKey(), time(), ['nx']);
} catch (\Throwable $e) {
// 记不下来只影响「卡多久放行」这一个判断,不该挡住刷库本身
}
}
/**
* 队首这批已经卡了多少秒,没卡住时返回 null
*
* @param Redis $redis
* @return int|null
*/
private static function stuckFor(Redis $redis): ?int
{
try {
$since = $redis->get(self::stuckSinceKey());
if ($since === false || !is_numeric($since)) {
return null;
}
return max(0, time() - (int)$since);
} catch (\Throwable $e) {
/*
* 读不到就当没卡住。宁可晚一点放行也不能早放行 ——
* 早放行等于把还能救的数据提前扔进死信。
*/
return null;
}
}
/**
* 批次确认掉了,卡住计时清零
*
* @param Redis $redis
* @return void
*/
private static function clearStuck(Redis $redis): void
{
try {
$redis->del(self::stuckSinceKey());
} catch (\Throwable $e) {
}
}
/**
* 锁的存活时间(秒),防止刷库进程挂掉后死锁
*
* 刷库期间会周期性续租,所以这个值不需要覆盖整次刷库的耗时,
* 只要能覆盖「两次续租之间」即可;取值越小,持锁进程被杀之后锁自然释放得越快。
*/
private const LOCK_TTL = 30;
/** 续租间隔(秒):刷库过程中每隔这么久把锁的存活时间顶回 LOCK_TTL */
private const LOCK_RENEW_INTERVAL = 10;
/** 队列长度硬上限,超出后丢弃最旧的记录,避免数据库长时间不可用时撑爆 Redis */
public const MAX_LENGTH = 200000;
/** 死信队列长度上限,超出后同样丢弃最旧的,避免脏数据把 Redis 撑爆 */
public const DEAD_MAX_LENGTH = 10000;
/**
* 队首批次卡住多久之后,把写不进去的行强行转进死信(秒)
*
* 写入失败按 WriteErrorKind 分类之后,除了明确的数据错,其余一律留着重试 ——
* 这是对的,但也意味着一条谁也认不出的错误可以永远占着队首,后面的消息
* 全部堵死,直到队列涨到 MAX_LENGTH 开始丢最旧的。那还是丢数据,只是换了个位置。
*
* 所以给「留着重试」加一个上界。取一整天是因为要盖过真实故障的修复时间:
* 磁盘满、权限配错、备库没切回来,这些通常几小时内有人处理;
* 取短了(比如几分钟)就等于把一次运维故障变成一次数据丢失,那正是要防的事。
*/
public const STUCK_SECONDS = 86400;
/** 单次刷库最多处理多少条,防止一次请求耗时过长 */
public const FLUSH_LIMIT = 5000;
/**
* 单次刷库的墙钟上限(秒)
*
* 条数上限挡不住「每条都很慢」的情况:数据库变慢、批量 INSERT 反复退化成逐行时,
* 5000 条也可能跑上好几分钟。锁虽然会续租,但一次刷库无限期占着队列本身就不健康。
* 这里再加一道时间闸门,超时就收工,剩下的留给下一次。
*/
public const FLUSH_DEADLINE = 20;
/**
* 各列允许的最大字符数,与 sql/*.sql 里的 varchar 长度一一对应
*
* 不截断的话,一条超长 UA 或 URL 会让整批 INSERT 失败,
* 然后退化成 1000 次逐行 INSERT —— 一条脏数据放大成一千次数据库往返。
* 而这个入口是匿名可达的。
*
* ip 的 39 单独说一句:这一列存的**不是地址文本**,而是地址的十进制整数表示
* (见 Core::ip62long())。上限来自 2^128-1 恰好是 39 位十进制数字,
* 和「完整展开的 IPv6 文本长 39 字符」只是碰巧同为 39,别按文本长度去推。
* 这里截断的后果也和别的列不同:其余列截掉尾巴只是信息变短,
* 而十进制数截掉末位等于除以 10 —— 存进去的是另一个看起来合法的地址。
*/
private const LIMITS = [
'ua' => 255,
'browser_id' => 32, 'browser_version' => 32,
'os_id' => 32, 'os_version' => 32,
'url' => 255, 'path' => 255, 'query_string' => 255,
'ip' => 39,
'entrypoint' => 255, 'entrypoint_domain' => 100,
'referer' => 255, 'referer_domain' => 100,
'robot_id' => 32, 'robot_version' => 32,
'event_id' => 32,
];
/** 这几列是 int unsigned,超出范围的值一律记为 null */
private const ID_COLUMNS = ['content_id', 'meta_id'];
/** int unsigned 的上界 */
private const UNSIGNED_INT_MAX = 4294967295;
/**
* 单条消息的字节上限
*
* normalize() 按列宽截断之后正常数据远到不了这个量级,
* 这道防线只为拦住结构本身就异常的输入。
*/
public const MAX_PAYLOAD_BYTES = 8192;
/** 入队字段,顺序固定;与 Migrate::COLUMNS 相同但不含自增主键 */
public const COLUMNS = [
'ua', 'browser_id', 'browser_version', 'os_id', 'os_version',
'url', 'path', 'query_string', 'ip', 'entrypoint', 'entrypoint_domain',
'referer', 'referer_domain', 'time', 'content_id', 'meta_id',
'robot', 'robot_id', 'robot_version',
'event_id',
];
/**
* 生成一条访问日志的唯一标识
*
* 队列做不到「恰好一次」:写库成功之后、从 processing 里确认之前进程被杀,
* 下一轮会把同一批再写一遍。有了这个标识,重复的那次会被唯一索引挡下来,
* 于是「至少一次」的投递变成了「恰好一次」的结果。
*
* 前 16 位十六进制是毫秒时间戳左移后补 16 位随机数,后 16 位纯随机:
* 时间在前使得标识按毫秒聚簇,唯一索引的写入集中在 B+ 树右端附近,
* 不会像纯随机标识那样每次插入都落到全表的随机页上。同一毫秒内仍是随机顺序。
*
* @return string 32 个十六进制字符
*/
public static function newEventId(): string
{
$ms = (int)(microtime(true) * 1000);
return bin2hex(pack('J', ($ms << 16) | random_int(0, 0xFFFF)))
. bin2hex(random_bytes(8));
}
/**
* 是否启用写入队列
* Redis 不可用时返回 false,调用方退回直写
*
* @param Redis|null $redis
* @param Config|array|null $config 插件配置
* @return bool
*/
public static function isEnabled(?Redis $redis, Config|array|null $config): bool
{
if ($redis === null) {
return false;
}
// 未显式关闭即为启用(Redis 已连上就说明用户配置过)
return !isset($config->writeQueue) || $config->writeQueue != '0';
}
/**
* 入队
*
* @param Redis $redis
* @param array $row
* @return bool 入队成功返回 true,失败由调用方退回直写
*/
public static function push(Redis $redis, array $row): bool
{
try {
$payload = json_encode(self::normalize($row), JSON_UNESCAPED_UNICODE);
if ($payload === false) {
return false;
}
/*
* 截断之后还能超限,说明这条记录的结构本身就不对(例如字段被塞成了数组)。
* 不入队,调用方会退回直写;直写同样会失败,但至少不会把异常数据
* 塞进队列去拖累整批 INSERT。
*/
if (strlen($payload) > self::MAX_PAYLOAD_BYTES) {
return false;
}
$length = $redis->rPush(self::key(), $payload);
if ($length === false) {
return false;
}
// 超出硬上限时丢掉最旧的部分,保留最新的 MAX_LENGTH 条
if ($length > self::MAX_LENGTH) {
$redis->lTrim(self::key(), -self::MAX_LENGTH, -1);
}
return true;
} catch (\Throwable $e) {
return false;
}
}
/**
* 队列长度,Redis 出错时按 0 处理
*
* 只在「拿不到长度就当没积压」无所谓的地方用(例如后台面板上的一个数字)。
* 需要区分「队列为空」和「Redis 故障」的调用方一律用 tryLength()。
*
* @param Redis $redis
* @return int
*/
public static function length(Redis $redis): int
{
return self::tryLength($redis) ?? 0;
}
/**
* 还没进数据库的消息条数,Redis 出错时返回 null
*
* 把故障伪装成 0 会让「队列为空」和「Redis 挂了」变成同一个返回值,
* 于是定时任务打印「队列为空,无需刷库」然后以成功退出,故障被彻底掩盖。
*
* 含 processing:那批已经离开主队列但还没落库,对「还剩多少没写」这个问题
* 它和主队列里的消息没有区别,漏算会让停用插件时误判成「已经刷干净了」。
*
* @param Redis $redis
* @return int|null
*/
public static function tryLength(Redis $redis): ?int
{
try {
$queue = $redis->lLen(self::key());
$processing = $redis->lLen(self::processingKey());
if ($queue === false || $processing === false) {
return null;
}
return (int)$queue + (int)$processing;
} catch (\Throwable $e) {
return null;
}
}
/**
* 正在写库(或上次没确认完)的条数
*
* 正常情况下要么是 0,要么是一个批次的大小;长期居高不下说明刷库一直在失败。
*
* @param Redis $redis
* @return int
*/
public static function processingLength(Redis $redis): int
{
try {
return (int)$redis->lLen(self::processingKey());
} catch (\Throwable $e) {
return 0;
}
}
/**
* 死信队列长度
*
* @param Redis $redis
* @return int
*/
public static function deadLength(Redis $redis): int
{
try {
return (int)$redis->lLen(self::deadKey());
} catch (\Throwable $e) {
return 0;
}
}
/**
* 原子地从主队列取走一批,并暂存到 processing
*
* 这里必须是一段 Lua,而不是 LRANGE + LTRIM 两条命令:
* 生产者在队列超过 MAX_LENGTH 时也会 LTRIM 裁掉队首,
* 一旦它插在消费者的「读」和「裁」之间,消费者按旧位置裁掉的
* 就是新写进来、还没落库的消息 —— 无声丢数据。
* Redis 执行脚本期间不处理别的命令,读和裁之间就没有缝可插了。
*
* 脚本内部的顺序是「先写 processing,后裁主队列」,不能反过来。
* Lua 脚本只保证不被别的命令打断,不保证出错回滚:先裁后写的话,
* 一旦 RPUSH 中途失败(processing 键类型不对等),消息就同时不在主队列
* 也不在 processing —— 无声丢失。反过来最坏只是主队列里留下副本,
* 下一轮重新取到,由 event_id 的唯一索引挡掉重复。
*
* @param Redis $redis
* @param int $count 最多取多少条
* @return array 取到的原始消息,顺序与队列一致
* @throws \RuntimeException 脚本执行失败时抛出,绝不能当成「队列已空」
*/
private static function claim(Redis $redis, int $count): array
{
$script = <<<'LUA'
local items = redis.call('LRANGE', KEYS[1], 0, ARGV[1] - 1)
if #items == 0 then
return items
end
for i = 1, #items do
redis.call('RPUSH', KEYS[2], items[i])
end
redis.call('LTRIM', KEYS[1], #items, -1)
return items
LUA;
$items = $redis->eval($script, [self::key(), self::processingKey(), $count], 2);
/*
* eval 出错时 phpredis 返回 false。以前这里 `is_array($items) ? $items : []`
* 把失败翻译成空数组,flush() 于是判定 stopped=empty 正常收工、退出码 0 ——
* 一条丢数据的路径连告警都没有。失败必须显式抛出。
*/
if ($items === false) {
$error = $redis->getLastError();
$redis->clearLastError();
throw new \RuntimeException('从队列取数失败:' . ($error !== null && $error !== '' ? $error : '未知错误'));
}
return is_array($items) ? $items : [];
}
/**
* 上一次没能确认的那批
*
* 消费者取走消息之后、写库成功之前挂掉,数据就停在 processing 里。
* 每轮开工前先把它捡回来,否则新取的一批会和它混在一起,没法分别确认。
*
* @param Redis $redis
* @return array
*/
private static function leftover(Redis $redis): array
{
$items = $redis->lRange(self::processingKey(), 0, -1);
/*
* lRange 失败返回 false,**绝不能当成「processing 是空的」**。
*
* 当成空的话流程会转去 claim() 取新一批,而 claim 的 Lua 是 RPUSH 到
* processing 尾部 —— 于是 processing 里前面压着上一批(还没落库),
* 后面接着新一批(马上要落库)。确认那一步是
* lTrim(processing, count($items), -1),按**位置**从头部裁掉 count 个,
* 砍掉的正是前面那批还没落库的消息,留下的反而是已经落库的。
* 砍错批次,直接丢数据。所以这里必须抛出去,让本轮刷库停手。
*/
if ($items === false) {
$error = $redis->getLastError();
$redis->clearLastError();
throw new \RuntimeException(
'读取 processing 失败:' . ($error !== null && $error !== '' ? $error : '未知错误')
);
}
return is_array($items) ? $items : [];
}
/**
* 把写不进去的消息挪进死信队列
*
* 必须在裁剪主队列之前调用:中途挂掉的话,宁可这批消息重复处理一次,
* 也不能出现「主队列已裁掉、死信里又没有」的空档。
*
* 写不进死信队列时抛异常而不是返回条数:调用方拿到条数之后会无条件
* 裁掉 processing,于是「没进死信、也没进数据库」的消息被无声删除 ——
* 这正是死信队列本身要防的事。抛出去让整批留在 processing 等下次重试。
*
* @param Redis $redis
* @param array $entries 每项为 ['reason' => string, 'payload' => string]
* @return int 实际入队条数
* @throws \RuntimeException 死信队列写入失败时抛出,调用方不得确认本批
*/
private static function pushDead(Redis $redis, array $entries): int
{
if (empty($entries)) {
return 0;
}
$payloads = [];
$at = time();
foreach ($entries as $entry) {
/*
* payload 是访客可控数据,可能不是合法 UTF-8,默认参数下 json_encode 会返回 false。
* 以前这里静默跳过,可这条消息随后照样被裁掉 —— 又是一条无声丢失。
* 先让非法字节被替换掉(内容仍可辨认),再退到 base64 保留原始字节。
*/
$encoded = json_encode([
'at' => $at,
'reason' => $entry['reason'],
'payload' => $entry['payload'],
], JSON_UNESCAPED_UNICODE | JSON_INVALID_UTF8_SUBSTITUTE);
if ($encoded === false) {
$encoded = json_encode([
'at' => $at,
'reason' => $entry['reason'] . '+unencodable',
'payload_b64' => base64_encode((string)$entry['payload']),
]);
}
if ($encoded === false) {
# 连 base64 都编码不出来说明结构本身异常,不能当作「已处理」放行
throw new \RuntimeException('死信队列消息无法编码,本批不予确认');
}
$payloads[] = $encoded;
}
$length = $redis->rPush(self::deadKey(), ...$payloads);
if ($length === false) {
throw new \RuntimeException(sprintf(
'死信队列写入失败(%d 条),本批保留在 processing 中等待重试',
count($payloads)
));
}
if ($length > self::DEAD_MAX_LENGTH) {
$redis->lTrim(self::deadKey(), -self::DEAD_MAX_LENGTH, -1);
}
return count($payloads);
}
/**
* 判断是否到了该刷库的时候
*
* @param Redis $redis
* @param int $size 条数阈值
* @param int $interval 时间阈值(秒)
* @return bool
*/
public static function isDue(Redis $redis, int $size, int $interval): bool
{
try {
$length = self::length($redis);
if ($length <= 0) {
return false;
}
if ($length >= max(1, $size)) {
return true;
}
$last = (int)$redis->get(self::lastFlushKey());
if ($last <= 0) {
// 没有记录过,先打上时间戳,下一轮再按间隔判断
$redis->set(self::lastFlushKey(), time());
return false;
}
return (time() - $last) >= max(1, $interval);
} catch (\Throwable $e) {
return false;
}
}
/**
* 抢刷库锁,拿到的请求才执行刷库
*
* 锁值是一次性随机 token 而不是 PID:PID 会复用,多机部署时更是会重复,
* 拿它当身份标识意味着「我」和「别人」根本分不开。
*
* @param Redis $redis
* @return string|null 抢到返回本次持有的 token,没抢到返回 null
*/
public static function acquireLock(Redis $redis): ?string
{
try {
$token = bin2hex(random_bytes(16));
$ok = $redis->set(self::lockKey(), $token, ['nx', 'ex' => self::LOCK_TTL]);
return $ok ? $token : null;
} catch (\Throwable $e) {
return null;
}
}
/**
* 续租:确认锁还在自己手上,然后把存活时间顶回 LOCK_TTL
*
* 返回 false 意味着锁已经不属于自己了(过期后被别人抢走,或被误删),
* 此时调用方必须立刻停手:继续刷下去就会和新的持有者同时读写、裁剪同一个队列。
*
* @param Redis $redis
* @param string $token acquireLock() 返回的 token
* @return bool
*/
public static function renewLock(Redis $redis, string $token): bool
{
try {
// 比较和续期必须是原子的,否则「比较通过 → 锁过期 → 续期」会把别人的锁续走
$script = <<<'LUA'
if redis.call("get", KEYS[1]) == ARGV[1] then
return redis.call("expire", KEYS[1], ARGV[2])
else
return 0
end
LUA;
return (int)$redis->eval($script, [self::lockKey(), $token, self::LOCK_TTL], 1) === 1;
} catch (\Throwable $e) {
return false;
}
}
/**
* 释放刷库锁
*
* 只删自己的锁:直接 DEL 的话,一个超时的旧消费者收尾时会把新消费者刚拿到的锁删掉,
* 于是第三个请求又能抢到锁,队列上同时出现多个消费者。
*
* @param Redis $redis
* @param string|null $token acquireLock() 返回的 token;null 表示没拿到锁,什么都不用做
* @return void
*/
public static function releaseLock(Redis $redis, ?string $token): void
{
if ($token === null) {
return;
}
try {
$script = <<<'LUA'
if redis.call("get", KEYS[1]) == ARGV[1] then
return redis.call("del", KEYS[1])
else
return 0
end
LUA;
$redis->eval($script, [self::lockKey(), $token], 1);
} catch (\Throwable $e) {
}
}
/**
* 把队列里的数据落库
*
* 先读后删:数据库不可用时保留队列不裁剪,等下次重试,宁可重复也不丢数据;
* 数据库是通的时候整批确认,写不进去的个别行先搬进死信队列再裁剪 ——
* 它们留在主队列里只会永远堵着,但直接丢掉就是无声的数据丢失。
*
* 上限按「从队列取走的条数」(attempted) 计算,而不是按写入成功数:
* 否则一批消息大部分写失败时 written 不增长,循环会继续吃下去,
* 极端情况(整个队列都是坏 JSON)下单次调用会一路处理完 MAX_LENGTH 条,
* 彻底失去时间边界。
*
* @param Redis $redis
* @param Db $db 统计数据所在的数据库
* @param int $limit 本次最多处理多少条,0 表示用默认上限
* @param float|null $deadline 墙钟截止时间(microtime 时间戳),null 表示用默认预算
* @param string|null $token 刷库锁的 token;传入后会在批次之间续租,一旦发现锁已易主立即停手
* @return array{attempted:int,written:int,invalid:int,rejected:int,dead:int,stopped:string,error:?string}
* attempted 取走并裁掉的消息数;written 实际入库行数;
* invalid JSON 解析失败的条数;rejected 数据库逐行重试仍写不进去的条数;
* dead 已转入死信队列的条数(invalid + rejected 都会进死信,不再无声丢弃);
* 恒有 invalid + rejected === attempted - written。
* stopped 为结束原因:
* empty 队列已取空(正常结束)
* limit 达到条数上限,队列中仍有积压
* deadline 达到时间上限,队列中仍有积压
* db 数据库不可用,本批已保留在队列中等待重试
* lock 刷库锁已经不在自己手上,为避免多消费者并发而主动停手;
* 未确认的那一批不计入 attempted/written,下次刷库会重放它
* error Redis 或其它异常(含死信写入失败、claim 脚本出错),详见 error;
* 未确认的批次同样留在 processing 里等下次重放
*/
public static function flush(
Redis $redis,
Db $db,
int $limit = 0,
?float $deadline = null,
?string $token = null
): array {
$limit = $limit > 0 ? $limit : self::FLUSH_LIMIT;
$deadline = $deadline ?? (microtime(true) + self::FLUSH_DEADLINE);
$batchSize = Migrate::BATCH_SIZE;
$renewAt = microtime(true) + self::LOCK_RENEW_INTERVAL;
$dates = []; // 本次写进数据库的记录覆盖了哪些日期,收尾时按它失效缓存
$result = [
'attempted' => 0,
'written' => 0,
'invalid' => 0,
'rejected' => 0,
'dead' => 0,
'invalidated' => 0,
'stopped' => 'empty',
'error' => null,
];
try {
while ($result['attempted'] < $limit) {
$now = microtime(true);
if ($now >= $deadline) {
$result['stopped'] = 'deadline';
break;
}
// 续租放在取数据之前:确认锁还在自己手上,再动队列
if ($token !== null && $now >= $renewAt) {
if (!self::renewLock($redis, $token)) {
$result['stopped'] = 'lock';
$result['error'] = '刷库锁已易主,为避免与另一个消费者同时裁剪队列而停止';
break;
}
$renewAt = $now + self::LOCK_RENEW_INTERVAL;
}
$take = min($batchSize, $limit - $result['attempted']);
# 先把上一轮没确认完的捡回来,没有残余才去主队列取新的
$items = self::leftover($redis);
$fromLeftover = !empty($items);
if (!$fromLeftover) {
$items = self::claim($redis, $take);
}
if (empty($items)) {
$result['stopped'] = 'empty';
break;
}
$rows = [];
$rowItem = []; // $rows 下标 => $items 下标,失败行要靠它找回原始消息
foreach ($items as $i => $item) {
$row = json_decode($item, true);
if (is_array($row)) {
$rowItem[] = $i;
$rows[] = self::normalize($row);
}
}
// insertBatchDetailed 不抛异常,返回写入行数、失败行下标和每个失败的归类
$outcome = empty($rows)
? ['written' => 0, 'failed' => [], 'kinds' => [], 'fatal' => null, 'error' => null]
: Migrate::insertBatchDetailed($db, $rows, self::COLUMNS);
$ok = $outcome['written'];
/*
* 语句还没发出去就整批失败(连不上、拼不出语句)—— 一行 failed 都拿不到,
* 没有任何依据说这批数据有问题,原样留着。
*/
$fatal = $outcome['fatal'];
if ($fatal !== null) {
/*
* 这条路也必须走 STUCK_SECONDS 那道闸门。
*
* 原来只 markStuck() 就 break,从不比较上限 —— 只有下面逐行失败那条路
* 会检查。于是一个持续性的连接层故障(库没了、账号被删、语句永远拼不出来)
* 会让这批永远停在 processing:主队列照常堆积,涨到 MAX_LENGTH 之后
* 从队首开始丢,丢的还是最早的那些。堵住的代价最终还是丢数据,只是换了个地方。
*/
$stuckFor = self::stuckFor($redis);
if ($stuckFor === null || $stuckFor < self::STUCK_SECONDS) {
self::markStuck($redis);
$result['stopped'] = 'db';
$result['error'] = sprintf(
'目标库写入失败(%s),本批 %d 条已保留在队列中等待重试,'
. '已卡住 %d 秒(上限 %d 秒):%s',
$fatal->name,
count($items),
(int)($stuckFor ?? 0),
self::STUCK_SECONDS,
(string)$outcome['error']
);
break;
}
/*
* 卡过头了,放行。这一批拿不到任何「哪几行有问题」的信息,
* 只能整批转死信 —— 至少留下原始内容可查、可回放,
* 而不是让它继续堵着直到主队列自己把更早的数据丢掉。
*/
$dead = [];
foreach ($items as $item) {
$dead[] = ['reason' => $fatal->reason(), 'payload' => $item];
}
$result['dead'] += self::pushDead($redis, $dead);
if ($redis->lTrim(self::processingKey(), count($items), -1) === false) {
$error = $redis->getLastError();
$redis->clearLastError();
throw new \RuntimeException(
'确认批次失败(LTRIM processing):'
. ($error !== null && $error !== '' ? $error : '未知错误')
);
}
self::clearStuck($redis);
$result['attempted'] += count($items);
$result['rejected'] += count($items);
$result['stopped'] = 'db';
$result['error'] = sprintf(
'目标库写入持续失败(%s)已达 %d 秒(上限 %d 秒),本批 %d 条转入死信队列以免堵死队列:%s',
$fatal->name,
(int)$stuckFor,
self::STUCK_SECONDS,
count($items),
(string)$outcome['error']
);
break;
}
/*
* 把失败行分成两堆:明确是这一行的错(转死信),和其余(留着重试)。
*
* 以前这里不分:只要 alive() 说数据库还活着,所有失败行一律 db-rejected
* 进死信。而 alive() 走的是 SELECT —— 实测库处于只读事务、或账号的
* INSERT 权限被回收时它照样成功,磁盘写满也一样。于是「每一条 INSERT
* 都失败」被判成「每一条都是脏数据」,整个队列被倒进死信,
* 死信满 DEAD_MAX_LENGTH 之后从最旧的开始丢 —— 静默的大规模数据丢失。
* 判据改成 SQLSTATE,见 Database::classifyWriteError() 与 WriteErrorKind。
*/
$rejected = []; // 下标 => true,明确的脏数据
$retry = []; // 下标 => WriteErrorKind,留着重试的
foreach ($outcome['failed'] as $i) {
$kind = $outcome['kinds'][$i] ?? WriteErrorKind::Unknown;
if ($kind->shouldRetry()) {
$retry[$i] = $kind;
} else {
$rejected[$i] = true;
}
}
/*
* 缓存失效的登记必须放在下面任何一个 break 之前:这些行已经躺在数据库里了,
* 后面无论因为什么原因没能确认这一批,缓存该失效的照样得失效。
* 数据变了而缓存没变,比刷库失败本身更难发现。
*/
if ($ok > 0) {
$failedRows = array_fill_keys($outcome['failed'], true);
$written = empty($outcome['failed'])
? $rows
: array_values(array_diff_key($rows, $failedRows));
foreach (Cache::datesOf($written) as $date) {
$dates[$date] = true;
}
}
/*
* 有行要留着重试就不能确认这一批:Redis List 只能按位置整批 LTRIM,
* 做不到「只确认成功的那几条」。整批留下,下轮重放 ——
* 已经写进去的那部分由 event_id 唯一索引挡掉重复。
*
* 唯一的例外是这批已经卡了太久(STUCK_SECONDS):那说明谁也没来修,
* 再留下去就是让它堵到队列涨满、从最旧的开始丢。到点了就放行,
* 把这些行按各自的归类转进死信,至少留下证据而不是无声消失。
*/
if (!empty($retry)) {
$stuckFor = self::stuckFor($redis);
if ($stuckFor === null || $stuckFor < self::STUCK_SECONDS) {
self::markStuck($redis);
$kinds = array_unique(array_map(static fn($k) => $k->name, $retry));
$result['stopped'] = 'db';
$result['error'] = sprintf(
'本批 %d 条中有 %d 条写入失败且判定为环境问题(%s),'
. '整批已保留在队列中等待重试,已卡住 %d 秒(上限 %d 秒):%s',
count($items),
count($retry),
implode('/', $kinds),
(int)($stuckFor ?? 0),
self::STUCK_SECONDS,
(string)$outcome['error']
);
break;
}
# 卡过头了,放行:这些行转死信,本批照常确认
$result['error'] = sprintf(
'本批 %d 条已卡住 %d 秒(超过 %d 秒上限),其中 %d 条写不进去的已转入死信队列',
count($items),
(int)$stuckFor,
self::STUCK_SECONDS,
count($retry)
);
}
$rejectedRows = $rejected + array_fill_keys(array_keys($retry), true);
/*
* 动 processing 之前最后一次确认锁还在自己手上。
*
* 循环顶部那次续租挡不住这个场景:整批 INSERT 失败会退化成
* 逐行写,一批 BATCH_SIZE 条就是上千次数据库往返,慢库上足以超过
* LOCK_TTL。锁一过期,另一个消费者拿到锁、读到同一个 processing,
* 两边再各自按 count($items) 做位置 LTRIM —— 后动手的那次砍掉的
* 就是对方刚取走、还没落库的消息。
*
* 停手的代价只是这批留在 processing 下轮重放,event_id 的唯一索引
* 会挡掉重复写入;继续往下走的代价是无声丢数据,两者不对等。
*/
if ($token !== null && !self::renewLock($redis, $token)) {
$result['stopped'] = 'lock';
$result['error'] = sprintf(
'写库期间刷库锁已易主,本批 %d 条(已写入 %d 行)保留在 processing 中未确认,'
. '下次刷库会重放,重复部分由 event_id 唯一索引挡下',
count($items),
$ok
);
break;
}
$renewAt = microtime(true) + self::LOCK_RENEW_INTERVAL;
/*
* 走到这里数据库是通的、锁也还在手上,这一批可以确认掉了。
* 但确认的前提是「没写进去的那些留下了证据」:Redis List 只能整批 LTRIM,
* 做不到只确认成功的几条,所以先把失败的原样搬进死信队列再裁剪。
* 顺序不能反 —— 反了就是老问题:一批 1000 条只成功 1 条,另外 999 条无声消失。
* pushDead() 写不进去会抛异常,异常会跳过下面的 LTRIM,这批原样留着。
*/
$rowOfItem = array_flip($rowItem);
$dead = [];
foreach ($items as $i => $item) {
if (!isset($rowOfItem[$i])) {
$dead[] = ['reason' => 'invalid-json', 'payload' => $item];
continue;
}
$row = $rowOfItem[$i];
if (!isset($rejectedRows[$row])) {
continue;
}
# 原因按归类记:db-rejected 是脏数据,db-environment/db-unknown 是卡过头才放行的
$kind = $retry[$row] ?? WriteErrorKind::Data;
$dead[] = ['reason' => $kind->reason(), 'payload' => $item];
}
$result['dead'] += self::pushDead($redis, $dead);
/*
* 确认:这批已经有归宿(进了数据库或死信),从 processing 清掉。
* lTrim 失败必须当场停手 —— 当成成功的话,这批会在下一轮被
* leftover() 再捡一次,而 attempted/written 已经按成功计过数了,
* 调用方(flush-queue.php 的退出码、后台提示)会读到一个假的成功。
*/
if ($redis->lTrim(self::processingKey(), count($items), -1) === false) {
$error = $redis->getLastError();
$redis->clearLastError();
throw new \RuntimeException(
'确认批次失败(LTRIM processing):'
. ($error !== null && $error !== '' ? $error : '未知错误')
);
}
# 这批确认掉了,卡住计时清零
self::clearStuck($redis);
$result['attempted'] += count($items);
$result['written'] += $ok;
$result['invalid'] += count($items) - count($rows);
$result['rejected'] += count($outcome['failed']);
# 残余批次的条数和 $take 无关,不能拿它推断队列空了
if (!$fromLeftover && count($items) < $take) {
// 队列里已经没有更多消息了
$result['stopped'] = 'empty';
break;
}
$result['stopped'] = 'limit';
}
$redis->set(self::lastFlushKey(), time());
} catch (\Throwable $e) {
// 刷库中断不影响调用方,未裁剪的数据仍在队列里
$result['stopped'] = 'error';
$result['error'] = $e->getMessage();
}
/*
* 放在 try 外面:中途出错时前面几批可能已经写进去了,那部分的缓存同样得失效。
* 数据已经变了而缓存还是旧的,比刷库失败本身更难发现。
*/