perf: 优化全文搜索

This commit is contained in:
kuaifan
2025-04-17 12:27:21 +08:00
parent 679c2070c1
commit f61e7caf2b
4 changed files with 817 additions and 115 deletions

View File

@@ -4,10 +4,9 @@ namespace App\Console\Commands;
use App\Models\WebSocketDialogMsg;
use App\Models\WebSocketDialogUser;
use App\Module\ElasticSearch\ElasticSearchKeyValue;
use App\Module\ElasticSearch\ElasticSearchUserMsg;
use App\Module\ZincSearch\ZincSearchKeyValue;
use App\Module\ZincSearch\ZincSearchUserMsg;
use Illuminate\Console\Command;
use Illuminate\Support\Facades\Log;
class SyncDialogUserMsgToZincSearch extends Command
{
@@ -22,7 +21,6 @@ class SyncDialogUserMsgToZincSearch extends Command
protected $signature = 'zinc:sync-dialog-user-msg {--f} {--i} {--c} {--batch=500}';
protected $description = '同步聊天会话用户和消息到 ZincSearch';
protected $client = null;
/**
* SyncDialogUserMsgToElasticsearch constructor.
@@ -30,12 +28,6 @@ class SyncDialogUserMsgToZincSearch extends Command
public function __construct()
{
parent::__construct();
try {
$this->es = new ElasticSearchUserMsg();
} catch (\Exception $e) {
$this->error('Elasticsearch连接失败: ' . $e->getMessage());
exit(1);
}
}
/**
@@ -44,37 +36,16 @@ class SyncDialogUserMsgToZincSearch extends Command
*/
public function handle()
{
$this->info('开始同步聊天数据...');
// 清除索引
if ($this->option('c')) {
$this->info('清除索引...');
if (!$this->es->indexExists()) {
$this->saveLastId(true);
$this->info('索引不存在');
return 0;
}
$result = $this->es->deleteIndex();
if (isset($result['error'])) {
$this->error('删除索引失败: ' . $result['error']);
return 1;
}
$this->saveLastId(true);
$this->info('索引删除成功');
ZincSearchUserMsg::clear();
ZincSearchKeyValue::clear();
$this->info("索引删除成功");
return 0;
}
// 判断创建索引
if (!$this->es->indexExists()) {
$this->info('创建索引...');
$result = ElasticSearchUserMsg::generateIndex();
if (isset($result['error'])) {
$this->error('创建索引失败: ' . $result['error']);
return 1;
}
$this->saveLastId(true);
$this->info('索引创建成功');
}
$this->info('开始同步聊天数据...');
// 同步用户-会话数据
$this->syncDialogUsers($this->option('batch'));
@@ -89,18 +60,12 @@ class SyncDialogUserMsgToZincSearch extends Command
/**
* 保存最后一个ID
* @param string|true $type
* @param integer $lastId
* @param string $type
* @param int $lastId
*/
private function saveLastId($type, $lastId = 0)
private function saveLastId(string $type, int $lastId = 0): void
{
if ($type === true) {
$setting = [];
} else {
$setting = ElasticSearchKeyValue::getArray('elasticSearch:sync');
$setting[$type] = $lastId;
}
ElasticSearchKeyValue::save('elasticSearch:sync', $setting);
ZincSearchKeyValue::set("sync", ["{$type}" => $lastId], true);
}
/**
@@ -108,11 +73,11 @@ class SyncDialogUserMsgToZincSearch extends Command
* @param $type
* @return int
*/
private function getLastId($type)
private function getLastId($type): int
{
if ($this->option('i')) {
$setting = ElasticSearchKeyValue::getArray('elasticSearch:sync');
return intval($setting[$type] ?? 0);
if ($this->option("i")) {
$array = ZincSearchKeyValue::getArray("sync");
return intval($array[$type] ?? 0);
}
return 0;
}
@@ -146,24 +111,7 @@ class SyncDialogUserMsgToZincSearch extends Command
$this->info("{$num}/{$count} ({$progress}%) 正在同步用户ID {$lastId} ~ {$dialogUsers->last()->id}");
// 批量索引数据
$params = ['body' => []];
foreach ($dialogUsers as $dialogUser) {
$params['body'][] = [
'index' => [
'_index' => ElasticSearchUserMsg::indexName(),
'_id' => ElasticSearchUserMsg::generateUserDicId($dialogUser),
]
];
$params['body'][] = ElasticSearchUserMsg::generateUserFormat($dialogUser);
}
if ($params['body']) {
$result = $this->es->bulk($params);
if (isset($result['errors']) && $result['errors']) {
$this->error('批量索引用户数据部分失败');
Log::error('Elasticsearch批量索引失败: ' . json_encode($result['items']));
}
}
ZincSearchUserMsg::batchSyncUsers($dialogUsers);
$lastId = $dialogUsers->last()->id;
$this->saveLastId('dialog_user', $lastId);
@@ -198,52 +146,8 @@ class SyncDialogUserMsgToZincSearch extends Command
$progress = round($num / $count * 100, 2);
$this->info("{$num}/{$count} ({$progress}%) 正在同步消息ID {$lastId} ~ {$dialogMsgs->last()->id}");
// 获取这些消息所属的会话对应的所有用户
$dialogIds = $dialogMsgs->pluck('dialog_id')->unique()->toArray();
$userDialogMap = [];
if (!empty($dialogIds)) {
$dialogUsers = WebSocketDialogUser::whereIn('dialog_id', $dialogIds)->get();
foreach ($dialogUsers as $dialogUser) {
$userDialogMap[$dialogUser->dialog_id][] = $dialogUser->userid;
}
}
// 批量索引消息数据
$params = ['body' => []];
foreach ($dialogMsgs as $dialogMsg) {
// 如果该会话没有用户,跳过
if (empty($userDialogMap[$dialogMsg->dialog_id])) {
continue;
}
// 为每个用户-会话关系创建子文档
foreach ($userDialogMap[$dialogMsg->dialog_id] as $userid) {
$params['body'][] = [
'index' => [
'_index' => ElasticSearchUserMsg::indexName(),
'_id' => ElasticSearchUserMsg::generateMsgDicId($dialogMsg, $userid),
'routing' => ElasticSearchUserMsg::generateMsgParentId($dialogMsg, $userid) // 路由到父文档
]
];
$params['body'][] = ElasticSearchUserMsg::generateMsgFormat($dialogMsg, $userid);
}
}
if (!empty($params['body'])) {
// 分批处理
$chunks = array_chunk($params['body'], 1000);
foreach ($chunks as $chunk) {
$chunkParams = ['body' => $chunk];
$result = $this->es->bulk($chunkParams);
if (isset($result['errors']) && $result['errors']) {
$this->error('批量索引消息数据部分失败');
Log::error('Elasticsearch批量索引失败: ' . json_encode($result['items']));
}
}
}
// 批量索引数据
ZincSearchUserMsg::batchSyncMsgs($dialogMsgs);
$lastId = $dialogMsgs->last()->id;
$this->saveLastId('dialog_msg', $lastId);