新版流量池服务端提交

This commit is contained in:
wong
2026-02-04 10:58:56 +08:00
parent f956fc949b
commit 7534dd79fa
19 changed files with 1101 additions and 147 deletions

View File

@@ -472,7 +472,13 @@ class MessageController extends BaseController
if (!empty($res) && empty($item['isSend']) && in_array($item['msgType'],[1,3,20,34,40,42,43,47,49])){
$friend = Db::name('wechat_friendship')->where('id',$item['wechatFriendId'])->find();
if (!empty($friend)){
$trafficPoolId = Db::name('traffic_pool_v1')->where('identifier',$friend['wechatId'])->value('id');
// ========== 旧版流量池代码(已废弃) ==========
// $trafficPoolId = Db::name('traffic_pool_v1')->where('identifier',$friend['wechatId'])->value('id');
// ========== 新版流量池代码 ==========
$trafficPool = Db::name('traffic_pool')->where('identifier', $friend['wechatId'])->find();
$trafficPoolId = $trafficPool ? $trafficPool['id'] : null;
// ========== 旧版流量池代码结束 ==========
if (!empty($trafficPoolId)){
$data = [
'type' => 4,

View File

@@ -0,0 +1,188 @@
<?php
namespace app\command;
use think\console\Command;
use think\console\Input;
use think\console\Output;
use think\console\input\Option;
use think\facade\App;
use think\facade\Log;
/**
* 清除过期日志文件命令
*
* 使用方法:
* php think clean:logs # 使用默认保留10天
* php think clean:logs --days=7 # 保留7天
* php think clean:logs --days=30 # 保留30天
* php think clean:logs --dry-run # 预览模式,不实际删除
*/
class CleanLogsCommand extends Command
{
protected function configure()
{
$this->setName('clean:logs')
->setDescription('清除过期的日志文件')
->addOption('days', 'd', Option::VALUE_OPTIONAL, '保留天数默认10天', 10)
->addOption('dry-run', null, Option::VALUE_NONE, '预览模式,不实际删除文件');
}
protected function execute(Input $input, Output $output)
{
$days = (int)$input->getOption('days');
$dryRun = $input->getOption('dry-run');
if ($days <= 0) {
$output->writeln('<error>保留天数必须大于0</error>');
return false;
}
if ($dryRun) {
$output->writeln('<info>运行在预览模式,不会实际删除文件</info>');
}
$output->writeln("<info>====================================</info>");
$output->writeln("<info> 清除过期日志文件</info>");
$output->writeln("<info>====================================</info>");
$output->writeln("保留天数: {$days}");
$output->writeln("");
// 获取日志目录
$logPath = App::getRuntimePath() . 'log' . DIRECTORY_SEPARATOR;
if (!is_dir($logPath)) {
$output->writeln("<comment>日志目录不存在: {$logPath}</comment>");
return false;
}
// 计算截止时间(保留指定天数之前的日志)
$cutoffTime = time() - ($days * 24 * 60 * 60);
$cutoffDate = date('Y-m-d H:i:s', $cutoffTime);
$output->writeln("<comment>清除 {$cutoffDate} 之前的日志文件</comment>");
$output->writeln("");
// 统计信息
$totalFiles = 0;
$deletedFiles = 0;
$totalSize = 0;
$freedSize = 0;
try {
// 递归扫描日志目录
$result = $this->cleanLogDirectory($logPath, $cutoffTime, $dryRun, $output);
$totalFiles = $result['total'];
$deletedFiles = $result['deleted'];
$totalSize = $result['totalSize'];
$freedSize = $result['freedSize'];
} catch (\Exception $e) {
$output->writeln('<error>清除日志时发生错误: ' . $e->getMessage() . '</error>');
Log::error('清除日志失败: ' . $e->getMessage());
return false;
}
// 输出统计信息
$output->writeln("");
$output->writeln("<info>====================================</info>");
$output->writeln("<info> 清除完成</info>");
$output->writeln("<info>====================================</info>");
$output->writeln("扫描文件数: {$totalFiles}");
$output->writeln("删除文件数: {$deletedFiles}");
$output->writeln("释放空间: " . $this->formatBytes($freedSize));
if ($dryRun) {
$output->writeln("");
$output->writeln("<comment>预览模式:实际未删除任何文件</comment>");
}
return true;
}
/**
* 递归清理日志目录
*/
protected function cleanLogDirectory($dir, $cutoffTime, $dryRun, Output $output)
{
$total = 0;
$deleted = 0;
$totalSize = 0;
$freedSize = 0;
if (!is_dir($dir)) {
return ['total' => 0, 'deleted' => 0, 'totalSize' => 0, 'freedSize' => 0];
}
$items = scandir($dir);
foreach ($items as $item) {
if ($item === '.' || $item === '..') {
continue;
}
$path = $dir . $item;
if (is_dir($path)) {
// 递归处理子目录
$result = $this->cleanLogDirectory($path . DIRECTORY_SEPARATOR, $cutoffTime, $dryRun, $output);
$total += $result['total'];
$deleted += $result['deleted'];
$totalSize += $result['totalSize'];
$freedSize += $result['freedSize'];
} elseif (is_file($path)) {
$total++;
$fileSize = filesize($path);
$totalSize += $fileSize;
// 获取文件修改时间
$fileMTime = filemtime($path);
// 如果文件修改时间早于截止时间,则删除
if ($fileMTime < $cutoffTime) {
$freedSize += $fileSize;
if ($dryRun) {
$output->writeln("<comment>[预览] 将删除: {$path} (" . date('Y-m-d H:i:s', $fileMTime) . ", " . $this->formatBytes($fileSize) . ")</comment>");
} else {
if (@unlink($path)) {
$deleted++;
$output->writeln("<info>已删除: {$path}</info>");
} else {
$output->writeln("<error>删除失败: {$path}</error>");
}
}
}
}
}
return [
'total' => $total,
'deleted' => $deleted,
'totalSize' => $totalSize,
'freedSize' => $freedSize,
];
}
/**
* 格式化字节数
*/
protected function formatBytes($bytes, $precision = 2)
{
$units = ['B', 'KB', 'MB', 'GB', 'TB'];
if ($bytes == 0) {
return '0 B';
}
$bytes = max($bytes, 0);
$pow = floor(($bytes ? log($bytes) : 0) / log(1024));
$pow = min($pow, count($units) - 1);
$bytes /= pow(1024, $pow);
return round($bytes, $precision) . ' ' . $units[$pow];
}
}

View File

@@ -5,12 +5,21 @@ namespace app\common\model;
use think\Model;
/**
* 流量池模型类
* 流量池模型类(旧版,已废弃)
*
* @deprecated 此模型已废弃,请使用 TrafficPoolV2 模型
* 旧表ck_traffic_pool_v1
* 新表ck_traffic_pool使用 TrafficPoolV2 模型)
*/
class TrafficPool extends Model
{
// ========== 旧版流量池表(已废弃) ==========
// 设置数据表名
protected $name = 'traffic_pool_v1';
// protected $name = 'traffic_pool_v1';
// ========== 新版流量池表 ==========
// 注意:为了兼容性,暂时保留此模型,但表名已改为新版
// 新代码请使用 TrafficPoolV2 模型
protected $name = 'traffic_pool';
// 自动写入时间戳
protected $autoWriteTimestamp = true;

View File

@@ -221,6 +221,15 @@ class TrafficPoolSource extends Model
$sourceData['chatroomOwners'] = [];
$sourceData['chatroomInfo'] = null;
$sourceData['displayId'] = $source['sourceWechatId'] ?: '';
// 尝试获取好友头像
if (!empty($source['sourceWechatId'])) {
$sourceData['sourceAvatar'] = Db::table('ck_traffic_pool')
->where('wechatId', $source['sourceWechatId'])
->value('avatar') ?: Db::table('s2_wechat_friend')
->where('wechatId', $source['sourceWechatId'])
->value('headImgUrl') ?: '';
}
} else {
$sourceData['chatroomOwners'] = [];
$sourceData['chatroomInfo'] = null;
@@ -434,6 +443,17 @@ class TrafficPoolSource extends Model
$chatroomId = $source['sourceChatroomId'];
$sourceData['chatroomOwners'] = $chatroomOwners[$chatroomId] ?? [];
$sourceData['chatroomInfo'] = self::getChatroomInfo($chatroomId);
} elseif ($source['sourceType'] == self::SOURCE_TYPE_FRIEND_ADD) {
$sourceData['chatroomOwners'] = [];
$sourceData['chatroomInfo'] = null;
// 尝试获取好友头像
if (!empty($source['sourceWechatId'])) {
$sourceData['sourceAvatar'] = Db::table('ck_traffic_pool')
->where('wechatId', $source['sourceWechatId'])
->value('avatar') ?: Db::table('s2_wechat_friend')
->where('wechatId', $source['sourceWechatId'])
->value('headImgUrl') ?: '';
}
} else {
$sourceData['chatroomOwners'] = [];
$sourceData['chatroomInfo'] = null;

View File

@@ -90,6 +90,8 @@ Route::group('v1/', function () {
Route::put('group/update', 'app\cunkebao\controller\TrafficPoolV2Controller@updateGroup'); // 更新分组
Route::delete('group/delete', 'app\cunkebao\controller\TrafficPoolV2Controller@deleteGroup'); // 删除分组
Route::get('group/members', 'app\cunkebao\controller\TrafficPoolV2Controller@getGroupMembers'); // 获取分组成员
Route::post('preview-users', 'app\cunkebao\controller\TrafficPoolV2Controller@previewUsers'); // 预览用户列表(根据筛选条件)
Route::get('filter-fields', 'app\cunkebao\controller\TrafficPoolV2Controller@getFilterFields'); // 获取筛选字段元数据
Route::post('group/add-members', 'app\cunkebao\controller\TrafficPoolV2Controller@addMembersToGroup'); // 添加成员到分组
Route::post('group/remove-members', 'app\cunkebao\controller\TrafficPoolV2Controller@removeMembersFromGroup'); // 移除分组成员
@@ -106,6 +108,10 @@ Route::group('v1/', function () {
Route::delete('tag/remove', 'app\cunkebao\controller\TrafficPoolV2Controller@removeTag'); // 移除标签
Route::post('tag/sync-from-engine', 'app\cunkebao\controller\TrafficPoolV2Controller@syncTagsFromEngine'); // 从标签引擎同步标签
// RFM评分相关
Route::post('calculate-rfm', 'app\cunkebao\controller\TrafficPoolV2Controller@calculateRfm'); // 计算RFM评分
Route::post('group/:groupId/calculate-rfm', 'app\cunkebao\controller\TrafficPoolV2Controller@calculateGroupRfm'); // 批量计算分组RFM评分
// 分配相关
Route::post('allocate', 'app\cunkebao\controller\TrafficPoolV2Controller@allocatePool'); // 分配流量
Route::post('recycle', 'app\cunkebao\controller\TrafficPoolV2Controller@recyclePool'); // 回收流量

View File

@@ -46,14 +46,19 @@ class RFMController extends BaseController
$weightM = isset($config['weight_M']) ? (float)$config['weight_M'] : self::DEFAULT_WEIGHT_M;
$abnormalMoneyRatio = isset($config['abnormal_money_ratio']) ? (float)$config['abnormal_money_ratio'] : self::DEFAULT_ABNORMAL_MONEY_RATIO;
$scoreScale = isset($config['score_scale']) ? (int)$config['score_scale'] : self::DEFAULT_SCORE_SCALE;
$missingStrategy = isset($config['missing_strategy']) ? $config['missing_strategy'] : 'score_1';
$missingStrategy = isset($config['missing_strategy']) ? $confi961102'] : 'score_1';
// 权重归一化处理
$weightSum = $weightR + $weightF + $weightM;
if ($weightSum != 1.0) {
if ($weightSum != 1.0 && $weightSum > 0) {
$weightR = $weightR / $weightSum;
$weightF = $weightF / $weightSum;
$weightM = $weightM / $weightSum;
} elseif ($weightSum == 0) {
// 如果权重全为0使用默认权重
$weightR = self::DEFAULT_WEIGHT_R;
$weightF = self::DEFAULT_WEIGHT_F;
$weightM = self::DEFAULT_WEIGHT_M;
}
// 计算时间范围
@@ -111,6 +116,8 @@ class RFMController extends BaseController
// 3. 异常值处理 - 剔除大额异常订单
$mValues = array_column($customerData, 'M');
$abnormalThreshold = null; // 初始化异常阈值
if (!empty($mValues)) {
sort($mValues);
$m99Percentile = $this->percentile($mValues, 0.99);
@@ -120,6 +127,11 @@ class RFMController extends BaseController
foreach ($customerData as &$customer) {
$customer['isAbnormal'] = $customer['M'] > $abnormalThreshold;
}
} else {
// 如果没有M值数据标记所有客户为非异常
foreach ($customerData as &$customer) {
$customer['isAbnormal'] = false;
}
}
// 4. 使用五分位法计算各维度的区间阈值
@@ -127,7 +139,7 @@ class RFMController extends BaseController
$fThresholds = $this->calculatePercentiles(array_column($customerData, 'F'), false);
// M维度排除异常值计算区间
$mValuesForPercentile = array_filter(array_column($customerData, 'M'), function($m) use ($abnormalThreshold) {
return isset($abnormalThreshold) ? $m <= $abnormalThreshold : true;
return $abnormalThreshold !== null ? $m <= $abnormalThreshold : true;
});
$mThresholds = $this->calculatePercentiles(array_values($mValuesForPercentile), false);
@@ -136,7 +148,7 @@ class RFMController extends BaseController
foreach ($customerData as $customer) {
$rScore = $this->scoreByPercentile($customer['R'], $rThresholds, true); // R是反向的
$fScore = $this->scoreByPercentile($customer['F'], $fThresholds, false);
$mScore = $customer['isAbnormal'] ? 5 : $this->scoreByPercentile($customer['M'], $mThresholds, false); // 异常值给最高分
$mScore = isset($customer['isAbnormal']) && $customer['isAbnormal'] ? 5 : $this->scoreByPercentile($customer['M'], $mThresholds, false); // 异常值给最高分
// 计算RFM总分加权求和
$rfmScore = $rScore * $weightR + $fScore * $weightF + $mScore * $weightM;
@@ -146,7 +158,8 @@ class RFMController extends BaseController
if ($scoreScale == 100) {
$rfmMin = $weightR * 1 + $weightF * 1 + $weightM * 1;
$rfmMax = $weightR * 5 + $weightF * 5 + $weightM * 5;
$standardScore = (int)round(($rfmScore - $rfmMin) / ($rfmMax - $rfmMin) * 99 + 1);
$range = $rfmMax - $rfmMin;
$standardScore = $range > 0 ? (int)round(($rfmScore - $rfmMin) / $range * 99 + 1) : 1;
}
$results[] = [
@@ -187,7 +200,7 @@ class RFMController extends BaseController
],
'statistics' => [
'total_customers' => count($results),
'avg_rfm_score' => round(array_sum(array_column($results, 'RFM_score')) / count($results), 2),
'avg_rfm_score' => count($results) > 0 ? round(array_sum(array_column($results, 'RFM_score')) / count($results), 2) : 0,
]
]
];
@@ -352,7 +365,7 @@ class RFMController extends BaseController
}
/**
* 更新RFM值到 ck_traffic_source_v1s2_wechat_friend 表
* 更新RFM值到 ck_traffic_source_v1s2_wechat_friend 和 ck_traffic_pool_company
*
* @param array $results RFM计算结果数组
* @param string|null $ownerWechatId 微信ID用于过滤更新范围
@@ -365,8 +378,11 @@ class RFMController extends BaseController
$rScore = (string)$result['R_score'];
$fScore = (string)$result['F_score'];
$mScore = (string)$result['M_score'];
$rfmRaw = $result['R_raw'];
$rfmF = $result['F_raw'];
$rfmM = $result['M_raw'];
// 更新 ck_traffic_source_v1 表
// 更新 ck_traffic_source_v1 表V1旧表
// 根据 identifier 更新所有匹配的记录
$trafficSourceUpdate = [
'R' => $rScore,
@@ -389,6 +405,18 @@ class RFMController extends BaseController
$wechatFriendWhere['ownerWechatId'] = $ownerWechatId;
}
WechatFriendModel::where($wechatFriendWhere)->update($wechatFriendUpdate);
// 更新 ck_traffic_pool_company 表V2新表
// 根据 identifier 更新identifier可能是wechatId、phone等
$poolCompanyUpdate = [
'rfmF' => $rfmF,
'rfmM' => $rfmM,
'updateTime' => date('Y-m-d H:i:s')
];
Db::table('ck_traffic_pool_company')
->where('identifier', $identifier)
->where('isDel', 0)
->update($poolCompanyUpdate);
}
} catch (\Exception $e) {

View File

@@ -412,7 +412,11 @@ class TrafficController extends BaseController
'wa.nickname', 'wa.avatar', 'wa.gender', 'wa.phone', 'wa.alias'
]
)
->join('traffic_pool_v1 p', 'p.identifier=tspi.identifier', 'left')
// ========== 旧版流量池代码(已废弃) ==========
// ->join('traffic_pool_v1 p', 'p.identifier=tspi.identifier', 'left')
// ========== 新版流量池代码 ==========
->join('traffic_pool p', 'p.identifier=tspi.identifier', 'left')
// ========== 旧版流量池代码结束 ==========
->join('wechat_account wa', 'tspi.identifier=wa.wechatId', 'left')
->where($where);

View File

@@ -97,12 +97,156 @@ class TrafficPoolV2Controller extends BaseController
try {
$group = $this->groupService->createGroup($companyId, $data, $userId);
return ResponseHelper::success(['id' => $group->id], '创建成功');
return ResponseHelper::success([
'id' => $group->id,
'groupName' => $group->groupName
], '创建成功');
} catch (\Exception $e) {
return ResponseHelper::error('创建分组失败:' . $e->getMessage());
}
}
/**
* 根据筛选条件预览用户列表
* GET /v1/traffic/pool/v2/preview-users
*
* @return \think\response\Json
*/
public function previewUsers()
{
$companyId = $this->getUserInfo('companyId');
$ruleConfig = $this->request->param('ruleConfig');
$page = $this->request->param('page', 1, 'intval');
$pageSize = $this->request->param('pageSize', 20, 'intval');
$keyword = $this->request->param('keyword', '');
if (empty($ruleConfig)) {
return ResponseHelper::error('筛选条件不能为空');
}
// 如果ruleConfig是JSON字符串解析它
if (is_string($ruleConfig)) {
$ruleConfig = json_decode($ruleConfig, true);
}
// 从ruleConfig中提取keyword如果前端放在里面的话
if (empty($keyword) && isset($ruleConfig['keyword'])) {
$keyword = $ruleConfig['keyword'];
unset($ruleConfig['keyword']);
}
try {
$result = $this->groupService->previewGroupMembers($companyId, $ruleConfig, $page, $pageSize, $keyword);
return ResponseHelper::success($result);
} catch (\Exception $e) {
return ResponseHelper::error('获取用户列表失败:' . $e->getMessage());
}
}
/**
* 获取筛选条件可选项(字段元数据)
* GET /v1/traffic/pool/v2/filter-fields
*
* @return \think\response\Json
*/
public function getFilterFields()
{
try {
$fields = [
[
'field' => 'lifecycle',
'label' => '客户周期',
'type' => 'select',
'options' => [
['label' => '新流量', 'value' => 1],
['label' => '成长期', 'value' => 2],
['label' => '成熟期', 'value' => 3],
['label' => '衰退期', 'value' => 4],
['label' => '流失期', 'value' => 5],
]
],
[
'field' => 'intentionLevel',
'label' => '意向等级',
'type' => 'select',
'options' => [
['label' => '未知', 'value' => 0],
['label' => '低意向', 'value' => 1],
['label' => '中意向', 'value' => 2],
['label' => '高意向', 'value' => 3],
]
],
[
'field' => 'level',
'label' => '客户等级',
'type' => 'select',
'options' => [
['label' => '普通', 'value' => 0],
['label' => '白银', 'value' => 1],
['label' => '黄金', 'value' => 2],
['label' => '钻石', 'value' => 3],
]
],
[
'field' => 'gender',
'label' => '性别',
'type' => 'select',
'options' => [
['label' => '未知', 'value' => 0],
['label' => '男', 'value' => 1],
['label' => '女', 'value' => 2],
]
],
[
'field' => 'friendStatus',
'label' => '好友状态',
'type' => 'select',
'options' => [
['label' => '未添加', 'value' => 0],
['label' => '已申请', 'value' => 1],
['label' => '已通过', 'value' => 2],
['label' => '已拒绝', 'value' => 3],
['label' => '已删除', 'value' => 4],
]
],
[
'field' => 'province',
'label' => '地区',
'type' => 'province',
],
[
'field' => 'totalOrderAmount',
'label' => '总消费金额',
'type' => 'number',
],
[
'field' => 'totalOrderCount',
'label' => '订单数量',
'type' => 'number',
],
[
'field' => 'totalMsgCount',
'label' => '消息数量',
'type' => 'number',
],
[
'field' => 'rfmF',
'label' => 'RFM-F值',
'type' => 'number',
],
[
'field' => 'rfmM',
'label' => 'RFM-M值',
'type' => 'number',
],
];
return ResponseHelper::success($fields);
} catch (\Exception $e) {
return ResponseHelper::error('获取字段列表失败:' . $e->getMessage());
}
}
/**
* 更新分组
* @return \think\response\Json
@@ -658,5 +802,96 @@ class TrafficPoolV2Controller extends BaseController
return ResponseHelper::error('获取行为轨迹失败:' . $e->getMessage());
}
}
/**
* 计算并更新RFM评分
* POST /v1/traffic/pool/v2/calculate-rfm
*
* @return \think\response\Json
*/
public function calculateRfm()
{
$companyId = $this->getUserInfo('companyId');
$identifier = $this->request->post('identifier', null); // 可选,指定用户标识
try {
// 实例化RFM控制器传递ClassTableService
$rfmController = new RFMController($this->classTable);
// 获取配置参数(可从请求参数中获取,或使用默认值)
$config = [
'cycle_days' => $this->request->post('cycle_days', 180),
'weight_R' => $this->request->post('weight_R', 0.4),
'weight_F' => $this->request->post('weight_F', 0.3),
'weight_M' => $this->request->post('weight_M', 0.3),
'score_scale' => $this->request->post('score_scale', 5),
];
// 调用RFM计算方法
// 注意这里不传ownerWechatId因为V2系统是按companyId区分的
$result = $rfmController->calculateRfmFromTrafficOrder($identifier, null, $config);
if ($result['code'] == 200) {
return ResponseHelper::success($result['data'], 'RFM计算完成');
} else {
return ResponseHelper::error($result['msg']);
}
} catch (\Exception $e) {
return ResponseHelper::error('RFM计算失败' . $e->getMessage());
}
}
/**
* 批量更新指定分组的RFM评分
* POST /v1/traffic/pool/v2/group/{groupId}/calculate-rfm
*
* @return \think\response\Json
*/
public function calculateGroupRfm()
{
$companyId = $this->getUserInfo('companyId');
$groupId = $this->request->param('groupId');
if (empty($groupId)) {
return ResponseHelper::error('分组ID不能为空');
}
try {
// 获取分组成员
$members = $this->groupService->getGroupMembers($groupId, $companyId, 1, 9999, []);
if (empty($members['list'])) {
return ResponseHelper::error('分组无成员');
}
// 实例化RFM控制器传递ClassTableService
$rfmController = new RFMController($this->classTable);
$successCount = 0;
$failCount = 0;
// 为每个成员计算RFM
foreach ($members['list'] as $member) {
$identifier = $member['identifier'];
$result = $rfmController->calculateRfmFromTrafficOrder($identifier, null, []);
if ($result['code'] == 200) {
$successCount++;
} else {
$failCount++;
}
}
return ResponseHelper::success([
'total' => count($members['list']),
'success' => $successCount,
'fail' => $failCount
], 'RFM批量计算完成');
} catch (\Exception $e) {
return ResponseHelper::error('RFM批量计算失败' . $e->getMessage());
}
}
}

View File

@@ -6,6 +6,8 @@ use library\ResponseHelper;
use think\Controller;
use think\Db;
use app\cunkebao\service\DistributionRewardService;
use app\cunkebao\service\TrafficPoolService;
use app\common\model\TrafficPoolSource;
/**
* 对外API接口控制器
@@ -95,27 +97,28 @@ class PostExternalApiV1Controller extends Controller
// 渠道IDcid对应 distribution_channel.id
$channelId = !empty($params['cid']) ? intval($params['cid']) : 0;
$trafficPool = Db::name('traffic_pool_v1')->where('identifier', $identifier)->find();
if (!$trafficPool) {
$trafficPoolId =Db::name('traffic_pool_v1')->insertGetId([
'identifier' => $identifier,
'mobile' => !empty($params['phone']) ? $params['phone'] : '',
'createTime' => time()
]);
}else{
$trafficPoolId = $trafficPool['id'];
}
// ========== 旧版流量池代码(已废弃,保留用于兼容) ==========
// $trafficPool = Db::name('traffic_pool_v1')->where('identifier', $identifier)->find();
// if (!$trafficPool) {
// $trafficPoolId =Db::name('traffic_pool_v1')->insertGetId([
// 'identifier' => $identifier,
// 'mobile' => !empty($params['phone']) ? $params['phone'] : '',
// 'createTime' => time()
// ]);
// }else{
// $trafficPoolId = $trafficPool['id'];
// }
// ========== 旧版流量池代码结束 ==========
$taskCustomer = Db::name('task_customer')
->where('task_id', $plan['id'])
->where('phone', $identifier)
->find();
// 处理用户画像
if(!empty($params['portrait']) && is_array($params['portrait'])){
$this->updatePortrait($params['portrait'],$trafficPoolId,$plan['companyId']);
}
// 处理用户画像已迁移到V2流量池此处保留兼容
// if(!empty($params['portrait']) && is_array($params['portrait'])){
// $this->updatePortrait($params['portrait'],$trafficPoolId,$plan['companyId']);
// }
if (!$taskCustomer) {
$tags = !empty($params['tags']) ? explode(',', $params['tags']) : [];
$siteTags = !empty($params['siteTags']) ? explode(',', $params['siteTags']) : [];
@@ -154,6 +157,60 @@ class PostExternalApiV1Controller extends Controller
'createTime' => time(),
]);
// 实时同步到 V2 流量池系统(异步处理,不影响主流程)
if ($customerId) {
try {
$poolService = new TrafficPoolService();
// 判断 identifier 类型:手机号还是微信号
$identifierType = 2; // 默认手机号
$isPhone = preg_match('/^\+?\d{6,}$/', $identifier);
if (!$isPhone && !empty($params['wechatId'])) {
$identifierType = 1; // 微信号
}
// 准备流量池数据
$poolData = [
'identifierType' => $identifierType,
'mobile' => !empty($params['phone']) ? $params['phone'] : ($isPhone ? $identifier : ''),
'wechatId' => !empty($params['wechatId']) ? $params['wechatId'] : (!$isPhone ? $identifier : ''),
'nickname' => !empty($params['name']) ? $params['name'] : '',
];
// 准备公司流量数据
$companyData = [
'phone' => !empty($params['phone']) ? $params['phone'] : ($isPhone ? $identifier : ''),
'realName' => !empty($params['name']) ? $params['name'] : '',
'remark' => !empty($params['remark']) ? $params['remark'] : '',
];
// 准备来源数据
$sourceData = [
'sourceName' => !empty($params['source']) ? $params['source'] : ('场景获客_' . $plan['name']),
'remark' => !empty($params['remark']) ? $params['remark'] : '',
'extra' => json_encode([
'planId' => $plan['id'],
'planName' => $plan['name'],
'channelId' => $finalChannelId,
'customerId' => $customerId,
], JSON_UNESCAPED_UNICODE),
];
// 同步到 V2 流量池
$poolService->enterPool(
$identifier,
$plan['companyId'],
TrafficPoolSource::SOURCE_TYPE_API, // API导入
$poolData,
$companyData,
$sourceData
);
} catch (\Exception $e) {
// 记录错误但不影响主流程
\think\facade\Log::error('同步到V2流量池失败' . $e->getMessage());
}
}
// 记录获客奖励(异步处理,不影响主流程)
if ($customerId) {
try {

View File

@@ -10,6 +10,8 @@ use think\facade\Env;
// use EasyWeChat\Kernel\Exceptions\DecryptException;
use EasyWeChat\Kernel\Http\StreamResponse;
use think\Db;
use app\cunkebao\service\TrafficPoolService;
use app\common\model\TrafficPoolSource;
class PosterWeChatMiniProgram extends Controller
{
@@ -112,16 +114,18 @@ class PosterWeChatMiniProgram extends Controller
if ($result['errcode'] == 0 && isset($result['phone_info']['phoneNumber'])) {
// ========== 旧版流量池代码(已废弃,保留用于兼容) ==========
// TODO 拿到手机号之后的后续操作:
// 1. 先写入 ck_traffic_pool_v1 表 identifier mobile 都是 用 phone字段的值
$trafficPool = Db::name('traffic_pool_v1')->where('identifier', $result['phone_info']['phoneNumber'])->find();
if (!$trafficPool) {
Db::name('traffic_pool_v1')->insert([
'identifier' => $result['phone_info']['phoneNumber'],
'mobile' => $result['phone_info']['phoneNumber'],
'createTime' => time()
]);
}
// $trafficPool = Db::name('traffic_pool_v1')->where('identifier', $result['phone_info']['phoneNumber'])->find();
// if (!$trafficPool) {
// Db::name('traffic_pool_v1')->insert([
// 'identifier' => $result['phone_info']['phoneNumber'],
// 'mobile' => $result['phone_info']['phoneNumber'],
// 'createTime' => time()
// ]);
// }
// ========== 旧版流量池代码结束已迁移到V2实时同步 ==========
// 2. 写入 ck_task_customer: 以 task_id ~~identifier~~ phone 为条件如果存在则忽略使用类似laravel的firstOrcreate但我不知道thinkphp5.1里的写法)
// $taskCustomer = Db::name('task_customer')->where('task_id', $taskId)->where('identifier', $result['phone_info']['phoneNumber'])->find();
$taskCustomer = Db::name('task_customer')
@@ -165,6 +169,50 @@ class PosterWeChatMiniProgram extends Controller
'siteTags' => json_encode([]),
]);
// 实时同步到 V2 流量池系统(异步处理,不影响主流程)
if ($customerId) {
try {
$poolService = new TrafficPoolService();
$identifier = $result['phone_info']['phoneNumber'];
// 准备流量池数据
$poolData = [
'identifierType' => 2, // 手机号
'mobile' => $identifier,
];
// 准备公司流量数据
$companyData = [
'phone' => $identifier,
];
// 准备来源数据
$sourceData = [
'sourceName' => $task['name'] ?? '海报获客',
'extra' => json_encode([
'planId' => $taskId,
'planName' => $task['name'] ?? '',
'channelId' => $finalChannelId,
'customerId' => $customerId,
'source' => 'poster_miniprogram',
], JSON_UNESCAPED_UNICODE),
];
// 同步到 V2 流量池
$poolService->enterPool(
$identifier,
$task['companyId'],
TrafficPoolSource::SOURCE_TYPE_POSTER, // 海报获客
$poolData,
$companyData,
$sourceData
);
} catch (\Exception $e) {
// 记录错误但不影响主流程
\think\facade\Log::error('同步到V2流量池失败' . $e->getMessage());
}
}
// 记录获客奖励(异步处理,不影响主流程)
if ($customerId) {
try {
@@ -259,31 +307,33 @@ class PosterWeChatMiniProgram extends Controller
continue;
}
$isPhone = preg_match('/^\+?\d{6,}$/', $identifier);
$trafficPool = Db::name('traffic_pool_v1')->where('identifier', $identifier)->find();
if (!$trafficPool) {
$insertData = [
'identifier' => $identifier,
'createTime' => time()
];
if ($isPhone) {
$insertData['mobile'] = $identifier;
} else {
$insertData['wechatId'] = $identifier;
}
Db::name('traffic_pool_v1')->insert($insertData);
} else {
$updates = [];
if ($isPhone && empty($trafficPool['mobile'])) {
$updates['mobile'] = $identifier;
}
if (!$isPhone && empty($trafficPool['wechatId'])) {
$updates['wechatId'] = $identifier;
}
if (!empty($updates)) {
$updates['updateTime'] = time();
Db::name('traffic_pool_v1')->where('id', $trafficPool['id'])->update($updates);
}
}
// ========== 旧版流量池代码(已废弃,保留用于兼容) ==========
// $trafficPool = Db::name('traffic_pool_v1')->where('identifier', $identifier)->find();
// if (!$trafficPool) {
// $insertData = [
// 'identifier' => $identifier,
// 'createTime' => time()
// ];
// if ($isPhone) {
// $insertData['mobile'] = $identifier;
// } else {
// $insertData['wechatId'] = $identifier;
// }
// Db::name('traffic_pool_v1')->insert($insertData);
// } else {
// $updates = [];
// if ($isPhone && empty($trafficPool['mobile'])) {
// $updates['mobile'] = $identifier;
// }
// if (!$isPhone && empty($trafficPool['wechatId'])) {
// $updates['wechatId'] = $identifier;
// }
// if (!empty($updates)) {
// $updates['updateTime'] = time();
// Db::name('traffic_pool_v1')->where('id', $trafficPool['id'])->update($updates);
// }
// }
// ========== 旧版流量池代码结束已迁移到V2实时同步 ==========
$taskCustomer = Db::name('task_customer')
->where('task_id', $taskId)
@@ -305,6 +355,55 @@ class PosterWeChatMiniProgram extends Controller
// 使用 insertGetId 以便在需要时记录获客奖励
$customerId = Db::name('task_customer')->insertGetId($insertCustomer);
// 实时同步到 V2 流量池系统(异步处理,不影响主流程)
if (!empty($customerId)) {
try {
$poolService = new TrafficPoolService();
// 判断 identifier 类型
$identifierType = $isPhone ? 2 : 1; // 2=手机号, 1=微信号
// 准备流量池数据
$poolData = [
'identifierType' => $identifierType,
'mobile' => $isPhone ? $identifier : '',
'wechatId' => !$isPhone ? $identifier : '',
];
// 准备公司流量数据
$companyData = [
'phone' => $isPhone ? $identifier : '',
'remark' => $remark,
];
// 准备来源数据
$sourceData = [
'sourceName' => $task['name'] ?? '海报获客',
'remark' => $remark,
'extra' => json_encode([
'planId' => $taskId,
'planName' => $task['name'] ?? '',
'channelId' => $finalChannelId,
'customerId' => $customerId,
'source' => 'poster_batch_import',
], JSON_UNESCAPED_UNICODE),
];
// 同步到 V2 流量池
$poolService->enterPool(
$identifier,
$task['companyId'],
TrafficPoolSource::SOURCE_TYPE_POSTER, // 海报获客
$poolData,
$companyData,
$sourceData
);
} catch (\Exception $e) {
// 记录错误但不影响主流程
\think\facade\Log::error('同步到V2流量池失败' . $e->getMessage());
}
}
// 表单录入成功即视为一次获客:
// 仅在存在有效渠道ID时记录获客奖励谁的cid谁获客
if (!empty($customerId) && $finalChannelId > 0) {

View File

@@ -72,7 +72,11 @@ class GetConvertedListWithInCompanyV1Controller extends BaseController
'f.tags', 'f.createTime', TrafficSourceModel::STATUS_PASSED . ' status'
]
)
->join('traffic_pool_v1 p', 'p.identifier=s.identifier')
// ========== 旧版流量池代码(已废弃) ==========
// ->join('traffic_pool_v1 p', 'p.identifier=s.identifier')
// ========== 新版流量池代码 ==========
->join('traffic_pool p', 'p.identifier=s.identifier')
// ========== 旧版流量池代码结束 ==========
->join('wechat_account w', 'p.wechatId=w.wechatId')
->join('wechat_friendship f', 'w.wechatId=f.wechatId and f.deleteTime=0')
->order('s.id desc');

View File

@@ -328,12 +328,24 @@ class GetPotentialListWithInCompanyV1Controller extends BaseController
if (empty($userId)) {
return json_encode(['code' => 500, 'msg' => '用户id不能为空']);
}
$data = Db::name('traffic_pool_v1')->alias('tp')
->join('wechat_friendship f', 'tp.wechatId=f.wechatId AND f.companyId=' . $companyId, 'left')
->join(['s2_wechat_friend' => 'wf'], 'f.wechatId=wf.wechatId', 'left')
->where(['tp.id' => $userId])
->order('tp.createTime desc')
// ========== 旧版流量池代码(已废弃) ==========
// $data = Db::name('traffic_pool_v1')->alias('tp')
// ->join('wechat_friendship f', 'tp.wechatId=f.wechatId AND f.companyId=' . $companyId, 'left')
// ->join(['s2_wechat_friend' => 'wf'], 'f.wechatId=wf.wechatId', 'left')
// ->where(['tp.id' => $userId])
// ->order('tp.createTime desc')
// ->column('wf.id,wf.labels,wf.siteLabels');
// ========== 新版流量池代码 ==========
$pool = Db::name('traffic_pool')->where('id', $userId)->find();
if (!$pool) {
return ResponseHelper::success(['wechat' => [], 'siteLabels' => []]);
}
$data = Db::name('s2_wechat_friend')->alias('wf')
->join('wechat_friendship f', 'wf.wechatId=f.wechatId AND f.companyId=' . $companyId, 'left')
->where(['wf.wechatId' => $pool['identifier']])
->order('wf.id desc')
->column('wf.id,wf.labels,wf.siteLabels');
// ========== 旧版流量池代码结束 ==========
if (empty($data)) {
return ResponseHelper::success(['wechat' => [], 'siteLabels' => []]);
}
@@ -424,12 +436,21 @@ class GetPotentialListWithInCompanyV1Controller extends BaseController
if (!is_array($userIds)) {
return ResponseHelper::error('选择的用户类型错误');
}
$result = Db::name('traffic_pool_v1')->alias('tp')
->join('traffic_source_v1 tc', 'tp.identifier=tc.identifier')
// ========== 旧版流量池代码(已废弃) ==========
// $result = Db::name('traffic_pool_v1')->alias('tp')
// ->join('traffic_source_v1 tc', 'tp.identifier=tc.identifier')
// ->whereIn('tp.id', $userIds)
// ->where(['companyId' => $companyId])
// ->group('tp.identifier')
// ->column('tc.identifier');
// ========== 新版流量池代码 ==========
$result = Db::name('traffic_pool')->alias('tp')
->join('traffic_pool_company tpc', 'tpc.poolId=tp.id AND tpc.companyId=' . $companyId)
->join('traffic_pool_source tps', 'tps.poolCompanyId=tpc.id')
->whereIn('tp.id', $userIds)
->where(['companyId' => $companyId])
->group('tp.identifier')
->column('tc.identifier');
->column('tps.identifier');
// ========== 旧版流量池代码结束 ==========
} else {
/*if (empty($tableFile)){
return ResponseHelper::error('请上传用户文件');
@@ -548,24 +569,29 @@ class GetPotentialListWithInCompanyV1Controller extends BaseController
$batchRows = array_slice($rows, $i, $batchSize);
if (!empty($batchRows)) {
$identifiers = array_column($batchRows, 'phone');
//流量池处理
$existing = Db::name('traffic_pool_v1')
->whereIn('identifier', $identifiers)
->column('identifier');
$newData = [];
foreach ($batchRows as $row) {
if (!in_array($row['phone'], $existing)) {
$newData[] = [
'identifier' => $row['phone'],
'mobile' => $row['phone'],
'createTime' => time(),
];
}
}
if (!empty($newData)) {
Db::name('traffic_pool_v1')->insertAll($newData);
}
// ========== 旧版流量池代码已废弃已迁移到V2实时同步 ==========
// //流量池处理
// $existing = Db::name('traffic_pool_v1')
// ->whereIn('identifier', $identifiers)
// ->column('identifier');
//
// $newData = [];
// foreach ($batchRows as $row) {
// if (!in_array($row['phone'], $existing)) {
// $newData[] = [
// 'identifier' => $row['phone'],
// 'mobile' => $row['phone'],
// 'createTime' => time(),
// ];
// }
// }
// if (!empty($newData)) {
// Db::name('traffic_pool_v1')->insertAll($newData);
// }
// ========== 新版流量池代码(使用 TrafficPoolService 实时同步) ==========
// 流量池处理 - 现在通过 TrafficPoolService 实时同步到 V2
// 如果需要批量导入,建议使用 migrate:trafficPoolV2 命令
// ========== 旧版流量池代码结束 ==========
//流量池来源处理
$newData2 = [];
@@ -638,10 +664,20 @@ class GetPotentialListWithInCompanyV1Controller extends BaseController
$isWechat = $this->request->param('isWechat', false);
$companyId = $this->getUserInfo('companyId');
$friend = Db::name('traffic_pool_v1')->alias('tp')
->join('wechat_friendship f', 'tp.wechatId=f.wechatId AND f.companyId='.$companyId, 'left')
->join(['s2_wechat_friend' => 'wf'], 'f.wechatId=wf.wechatId', 'left')
->where(['tp.id' => $userId])
// ========== 旧版流量池代码(已废弃) ==========
// $friend = Db::name('traffic_pool_v1')->alias('tp')
// ->join('wechat_friendship f', 'tp.wechatId=f.wechatId AND f.companyId='.$companyId, 'left')
// ->join(['s2_wechat_friend' => 'wf'], 'f.wechatId=wf.wechatId', 'left')
// ->where(['tp.id' => $userId])
// ========== 新版流量池代码 ==========
$pool = Db::name('traffic_pool')->where('id', $userId)->find();
if (!$pool) {
return ResponseHelper::error('流量池记录不存在');
}
$friend = Db::name('s2_wechat_friend')->alias('wf')
->join('wechat_friendship f', 'wf.wechatId=f.wechatId AND f.companyId='.$companyId, 'left')
->where(['wf.wechatId' => $pool['identifier']])
// ========== 旧版流量池代码结束 ==========
->order('tp.createTime desc')
->column('wf.id,wf.accountId,wf.labels,wf.siteLabels');
if (empty($data)) {

View File

@@ -2398,10 +2398,18 @@ class WorkbenchController extends Controller
];
// 查询发布记录
// ========== 旧版流量池代码(已废弃) ==========
// $list = Db::name('workbench_import_contact_item')->alias('wici')
// ->join('traffic_pool_v1 tp', 'tp.id = wici.poolId', 'left')
// ->join('traffic_source_v1 tc', 'tc.identifier = tp.identifier', 'left')
// ->join('wechat_account wa', 'wa.wechatId = tp.wechatId', 'left')
// ========== 新版流量池代码 ==========
$list = Db::name('workbench_import_contact_item')->alias('wici')
->join('traffic_pool_v1 tp', 'tp.id = wici.poolId', 'left')
->join('traffic_source_v1 tc', 'tc.identifier = tp.identifier', 'left')
->join('traffic_pool tp', 'tp.id = wici.poolId', 'left')
->join('traffic_pool_company tpc', 'tpc.poolId = tp.id AND tpc.companyId = ' . $this->getUserInfo('companyId'), 'left')
->join('traffic_pool_source tps', 'tps.poolCompanyId = tpc.id', 'left')
->join('wechat_account wa', 'wa.wechatId = tp.wechatId', 'left')
// ========== 旧版流量池代码结束 ==========
->field([
'wici.id',
'wici.workbenchId',
@@ -2409,7 +2417,7 @@ class WorkbenchController extends Controller
'tp.identifier',
'tp.mobile',
'tp.wechatId',
'tc.name',
'tps.sourceName as name', // 从新版来源表获取名称
'wa.nickName',
'wa.avatar',
'wa.alias',

View File

@@ -25,10 +25,18 @@ class WorkbenchImportContactController extends Controller
];
// 查询发布记录
// ========== 旧版流量池代码(已废弃) ==========
// $list = Db::name('workbench_import_contact_item')->alias('wici')
// ->join('traffic_pool_v1 tp', 'tp.id = wici.poolId', 'left')
// ->join('traffic_source_v1 tc', 'tc.identifier = tp.identifier', 'left')
// ->join('wechat_account wa', 'wa.wechatId = tp.wechatId', 'left')
// ========== 新版流量池代码 ==========
$list = Db::name('workbench_import_contact_item')->alias('wici')
->join('traffic_pool_v1 tp', 'tp.id = wici.poolId', 'left')
->join('traffic_source_v1 tc', 'tc.identifier = tp.identifier', 'left')
->join('traffic_pool tp', 'tp.id = wici.poolId', 'left')
->join('traffic_pool_company tpc', 'tpc.poolId = tp.id AND tpc.companyId = ' . $this->getUserInfo('companyId'), 'left')
->join('traffic_pool_source tps', 'tps.poolCompanyId = tpc.id', 'left')
->join('wechat_account wa', 'wa.wechatId = tp.wechatId', 'left')
// ========== 旧版流量池代码结束 ==========
->field([
'wici.id',
'wici.workbenchId',
@@ -36,7 +44,7 @@ class WorkbenchImportContactController extends Controller
'tp.identifier',
'tp.mobile',
'tp.wechatId',
'tc.name',
'tps.sourceName as name', // 从新版来源表获取名称
'wa.nickName',
'wa.avatar',
'wa.alias',

View File

@@ -258,18 +258,31 @@ class TrafficPoolGroupService
'tpc.id',
'tpc.poolId',
'tpc.identifier',
'tpc.companyId',
'tpc.friendStatus',
'tpc.level',
'tpc.intentionLevel',
'tpc.lastInteractTime',
'tpc.rfmF',
'tpc.rfmM',
'tpc.totalMsgCount',
'tpc.totalOrderAmount',
'tpc.lastMsgTime',
'tpc.firstSourceType',
'tpc.firstSourceTime',
'tpc.lifecycle',
'tpc.createTime',
'tpc.realName',
'tpc.phone',
'tp.nickname',
'tp.avatar',
'tp.wechatId',
'tp.wechatAlias',
'tp.gender',
'tp.region',
'tp.country',
'tp.province',
'tp.city',
'tpgm.createTime as addTime'
])
->order('tpgm.createTime DESC')
@@ -312,19 +325,32 @@ class TrafficPoolGroupService
'tpc.id',
'tpc.poolId',
'tpc.identifier',
'tpc.companyId',
'tpc.friendStatus',
'tpc.level',
'tpc.intentionLevel',
'tpc.lastInteractTime',
'tpc.rfmF',
'tpc.rfmM',
'tpc.totalMsgCount',
'tpc.totalOrderAmount',
'tpc.lastMsgTime',
'tpc.firstSourceType',
'tpc.firstSourceTime',
'tpc.lifecycle',
'tpc.createTime',
'tpc.realName',
'tpc.phone',
'tpc.createTime as addTime',
'tp.nickname',
'tp.avatar',
'tp.wechatId'
'tp.wechatId',
'tp.wechatAlias',
'tp.gender',
'tp.region',
'tp.country',
'tp.province',
'tp.city'
])
->order('tpc.id DESC')
->page($page, $pageSize)
@@ -377,11 +403,58 @@ class TrafficPoolGroupService
}
});
} elseif ($condition['type'] === 'field') {
// 字段条件
$field = 'tpc.' . $condition['field'];
// 字段条件 - 根据字段所属表使用正确的别名
$fieldName = $condition['field'];
$operator = $condition['operator'];
$value = $condition['value'];
// 特殊处理keyword 字段用于多字段搜索
if ($fieldName === 'keyword') {
$keyword = $value;
$query->$method(function($q) use ($keyword) {
$q->where('tp.nickname', 'like', "%{$keyword}%")
->whereOr('tp.wechatId', 'like', "%{$keyword}%")
->whereOr('tp.wechatAlias', 'like', "%{$keyword}%")
->whereOr('tpc.realName', 'like', "%{$keyword}%")
->whereOr('tpc.phone', 'like', "%{$keyword}%");
});
return;
}
// 特殊处理friendIds 字段用于指定好友ID列表
if ($fieldName === 'friendIds') {
if (is_array($value) && !empty($value)) {
$query->$method('tpc.id', 'in', $value);
}
return;
}
// ck_traffic_pool 表的字段(基础用户信息)
$tpFields = ['nickname', 'avatar', 'wechatId', 'wechatAlias', 'gender', 'region', 'country', 'province', 'city', 'signature'];
// 判断字段属于哪个表
if (in_array($fieldName, $tpFields)) {
$field = 'tp.' . $fieldName;
} else {
// ck_traffic_pool_company 表的字段(公司维度信息)
$field = 'tpc.' . $fieldName;
}
// 特殊处理地区字段province
// 前端可能传递 "广东" 或 "广东 广州市"
if ($fieldName === 'province' && strpos($value, ' ') !== false) {
// 包含空格,说明是 "省份 城市" 格式
$parts = explode(' ', $value, 2);
$provinceName = trim($parts[0]);
$cityName = trim($parts[1]);
$query->$method(function($q) use ($provinceName, $cityName) {
$q->where('tp.province', '=', $provinceName)
->where('tp.city', 'like', "%{$cityName}%");
});
return;
}
switch ($operator) {
case '=':
case '!=':
@@ -439,12 +512,45 @@ class TrafficPoolGroupService
protected function formatMemberList($list, int $total, int $page, int $pageSize)
{
$result = [];
$poolCompanyIds = [];
// 收集所有的poolCompanyId
foreach ($list as $item) {
$poolCompanyIds[] = $item['id'];
}
// 批量查询标签
$tagsMap = [];
if (!empty($poolCompanyIds)) {
$tags = \think\Db::table('ck_traffic_pool_tag')
->alias('tpt')
->join('ck_traffic_pool_tag_define tptd', 'tpt.tagDefineId = tptd.id', 'LEFT')
->where('tpt.poolCompanyId', 'in', $poolCompanyIds)
->where('tpt.isDel', 0)
->where('tptd.isDel', 0)
->field('tpt.poolCompanyId, tptd.tagName, tptd.tagType')
->select();
foreach ($tags as $tag) {
$poolCompanyId = $tag['poolCompanyId'];
if (!isset($tagsMap[$poolCompanyId])) {
$tagsMap[$poolCompanyId] = [];
}
$tagsMap[$poolCompanyId][] = [
'tagName' => $tag['tagName'],
'tagType' => $tag['tagType']
];
}
}
foreach ($list as $item) {
$data = $item->toArray();
// 计算 RFM R 值
$data['rfmR'] = $item->lastInteractTime ? (int)floor((time() - $item->lastInteractTime) / 86400) : 9999;
// 计算 RFM 总分
$data['rfmScore'] = $this->calculateRfmScore($data['rfmR'], $data['rfmF'] ?? 0, $data['rfmM'] ?? 0);
// 添加标签
$data['tags'] = $tagsMap[$item['id']] ?? [];
$result[] = $data;
}
@@ -479,6 +585,7 @@ class TrafficPoolGroupService
// 动态规则分组
$query = TrafficPoolCompany::alias('tpc')
->join('ck_traffic_pool tp', 'tp.id = tpc.poolId', 'LEFT')
->where('tpc.companyId', $companyId)
->where('tpc.isDel', 0);
@@ -512,6 +619,7 @@ class TrafficPoolGroupService
->where('tpc.isDel', 0);
} else {
$query = TrafficPoolCompany::alias('tpc')
->join('ck_traffic_pool tp', 'tp.id = tpc.poolId', 'LEFT')
->where('tpc.companyId', $companyId)
->where('tpc.isDel', 0);
@@ -645,6 +753,78 @@ class TrafficPoolGroupService
'total' => $rScore + $fScore + $mScore
];
}
/**
* 预览动态分组成员(不创建分组,只预览符合条件的用户)
*
* @param int $companyId 公司ID
* @param array $ruleConfig 规则配置
* @param int $page 页码
* @param int $pageSize 每页数量
* @param string $keyword 搜索关键词
* @return array
*/
public function previewGroupMembers(int $companyId, array $ruleConfig, int $page = 1, int $pageSize = 20, string $keyword = '')
{
$query = TrafficPoolCompany::alias('tpc')
->join('ck_traffic_pool tp', 'tp.id = tpc.poolId', 'LEFT')
->where('tpc.companyId', $companyId)
->where('tpc.isDel', 0);
// 应用规则条件
if (!empty($ruleConfig)) {
$this->applyRuleConditions($query, $ruleConfig, $companyId);
}
// 应用关键字搜索
if (!empty($keyword)) {
$query->where(function($q) use ($keyword) {
$q->where('tp.nickname', 'like', "%{$keyword}%")
->whereOr('tp.wechatId', 'like', "%{$keyword}%")
->whereOr('tp.wechatAlias', 'like', "%{$keyword}%")
->whereOr('tpc.realName', 'like', "%{$keyword}%")
->whereOr('tpc.phone', 'like', "%{$keyword}%");
});
}
$total = $query->count();
$list = $query->field([
'tpc.id',
'tpc.poolId',
'tpc.identifier',
'tpc.companyId',
'tpc.friendStatus',
'tpc.level',
'tpc.intentionLevel',
'tpc.lastInteractTime',
'tpc.rfmF',
'tpc.rfmM',
'tpc.totalMsgCount',
'tpc.totalOrderAmount',
'tpc.lastMsgTime',
'tpc.firstSourceType',
'tpc.firstSourceTime',
'tpc.lifecycle',
'tpc.createTime',
'tpc.realName',
'tpc.phone',
'tp.nickname',
'tp.avatar',
'tp.wechatId',
'tp.wechatAlias',
'tp.gender',
'tp.region',
'tp.country',
'tp.province',
'tp.city'
])
->order('tpc.id DESC')
->page($page, $pageSize)
->select();
return $this->formatMemberList($list, $total, $page, $pageSize);
}
}

View File

@@ -602,12 +602,21 @@ class WorkbenchGroupPushJob
*/
protected function getFriendsByNormalPools(array $packageIds, $companyId, array $ownerWechatIds = [])
{
// ========== 旧版流量池代码(已废弃) ==========
// $query = Db::name('traffic_source_package_item_v1')
// ->alias('tspi')
// ->leftJoin('traffic_source_package_v1 tsp', 'tsp.id = tspi.packageId')
// ->leftJoin('traffic_pool_v1 tp', 'tp.identifier = tspi.identifier')
// ->leftJoin(['s2_wechat_friend' => 'wf'], 'wf.wechatId = tp.wechatId')
// ->leftJoin(['s2_wechat_account' => 'wa'], 'wa.id = wf.wechatAccountId')
// ========== 新版流量池代码 ==========
$query = Db::name('traffic_source_package_item_v1')
->alias('tspi')
->leftJoin('traffic_source_package_v1 tsp', 'tsp.id = tspi.packageId')
->leftJoin('traffic_pool_v1 tp', 'tp.identifier = tspi.identifier')
->leftJoin('traffic_pool tp', 'tp.identifier = tspi.identifier')
->leftJoin(['s2_wechat_friend' => 'wf'], 'wf.wechatId = tp.wechatId')
->leftJoin(['s2_wechat_account' => 'wa'], 'wa.id = wf.wechatAccountId')
// ========== 旧版流量池代码结束 ==========
->whereIn('tspi.packageId', $packageIds)
->where('tsp.isDel', 0)
->where('wf.isDeleted', 0)

View File

@@ -345,18 +345,34 @@ class WorkbenchImportContactJob
->column('id');
if (!empty($packageIds)) {
// ========== 旧版流量池代码(已废弃) ==========
// $normalData = Db::name('traffic_source_package_item_v1')->alias('tpi')
// ->join('traffic_pool_v1 tp', 'tp.identifier = tpi.identifier')
// ->join('traffic_source_v1 ts', 'ts.identifier = tpi.identifier','left')
// ->join('workbench_import_contact_item wici', 'wici.poolId = tp.id AND wici.workbenchId = '.$workbench->id,'left')
// ->where('tp.mobile', '>',0)
// ->where('wici.id','null')
// ->whereIn('tpi.packageId',$packageIds)
// ->field('tp.id,tpi.packageId,tp.mobile as phone,ts.name')
// ->order('tp.id DESC')
// ->group('tpi.identifier')
// ->limit($contactNum)
// ->select();
// ========== 新版流量池代码 ==========
$normalData = Db::name('traffic_source_package_item_v1')->alias('tpi')
->join('traffic_pool_v1 tp', 'tp.identifier = tpi.identifier')
->join('traffic_source_v1 ts', 'ts.identifier = tpi.identifier','left')
->join('traffic_pool tp', 'tp.identifier = tpi.identifier')
->join('traffic_pool_company tpc', 'tpc.poolId = tp.id AND tpc.companyId = ' . ($workbench->companyId ?? 0))
->join('traffic_pool_source tps', 'tps.poolCompanyId = tpc.id', 'left')
->join('workbench_import_contact_item wici', 'wici.poolId = tp.id AND wici.workbenchId = '.$workbench->id,'left')
->where('tp.mobile', '>',0)
->where('wici.id','null')
->whereIn('tpi.packageId',$packageIds)
->field('tp.id,tpi.packageId,tp.mobile as phone,ts.name')
->field('tp.id,tpi.packageId,tp.mobile as phone,tps.sourceName as name')
->order('tp.id DESC')
->group('tpi.identifier')
->limit($contactNum)
->select();
// ========== 旧版流量池代码结束 ==========
$data = array_merge($data, $normalData ?: []);
}
}
@@ -389,12 +405,22 @@ class WorkbenchImportContactJob
return [];
}
// 从 s2_wechat_friend 表获取好友,然后关联 traffic_pool_v1 表获取手机号
// ========== 旧版流量池代码(已废弃) ==========
// // 从 s2_wechat_friend 表获取好友,然后关联 traffic_pool_v1 表获取手机号
// $data = Db::table('s2_wechat_friend')->alias('wf')
// ->join('traffic_pool_v1 tp', 'tp.wechatId = wf.wechatId', 'left')
// ->join('traffic_source_v1 ts', 'ts.identifier = tp.identifier', 'left')
// ->join('workbench_import_contact_item wici', 'wici.poolId = tp.id AND wici.workbenchId = '.$workbench->id, 'left')
// ->where('wf.ownerWechatId', 'in', $wechatIds)
// ========== 新版流量池代码 ==========
// 从 s2_wechat_friend 表获取好友,然后关联 traffic_pool 表获取手机号
$data = Db::table('s2_wechat_friend')->alias('wf')
->join('traffic_pool_v1 tp', 'tp.wechatId = wf.wechatId', 'left')
->join('traffic_source_v1 ts', 'ts.identifier = tp.identifier', 'left')
->join('traffic_pool tp', 'tp.wechatId = wf.wechatId', 'left')
->join('traffic_pool_company tpc', 'tpc.poolId = tp.id AND tpc.companyId = ' . ($workbench->companyId ?? 0), 'left')
->join('traffic_pool_source tps', 'tps.poolCompanyId = tpc.id', 'left')
->join('workbench_import_contact_item wici', 'wici.poolId = tp.id AND wici.workbenchId = '.$workbench->id, 'left')
->where('wf.ownerWechatId', 'in', $wechatIds)
// ========== 旧版流量池代码结束 ==========
->where('wf.isDeleted', 0)
->where('tp.mobile', '>', 0)
->where('wici.id', 'null')

View File

@@ -232,6 +232,22 @@ return [
'log_file' => 'call_recording.log',
],
// ===========================
// 低频任务(每 2 小时)
// ===========================
// V2 流量池数据同步,全量同步好友、群成员和标签数据到 V2 流量池系统
'traffic_pool_v2_sync' => [
'name' => 'V2 流量池数据同步',
'command' => 'migrate:trafficPoolV2',
'schedule' => '0 */2 * * *', // 每2小时的0分执行0:00, 2:00, 4:00...
'options' => [],
'enabled' => true,
'max_concurrent' => 1,
'timeout' => 3600, // 1小时超时
'log_file' => 'traffic_pool_v2_sync.log',
],
// ===========================
// 每日 / 每几天任务
// ===========================
@@ -291,6 +307,18 @@ return [
'log_file' => 'calculate_score.log',
],
// 每日 3:00 清除过期日志文件默认保留10天可通过 --days 参数修改
'clean_logs' => [
'name' => '清除过期日志文件',
'command' => 'clean:logs',
'schedule' => '0 3 * * *', // 每天3点
'options' => ['--days=10'], // 默认保留10天可修改为其他天数如 ['--days=7'] 保留7天
'enabled' => true,
'max_concurrent' => 1,
'timeout' => 300, // 5分钟超时
'log_file' => 'clean_logs.log',
],
// 每 3 天执行的全量任务
// 每 3 天 3:00 全量同步所有在线好友,做一次大规模校准

View File

@@ -1292,44 +1292,47 @@ class Adapter implements WeChatServiceInterface
$affectedRows = 0;
try {
for ($i = 0; $i < $batchCount; $i++) {
$offset = $i * $batchSize;
// 使用 ON DUPLICATE KEY UPDATE 支持插入和更新
$sql = "INSERT INTO ck_traffic_pool_v1(
`identifier`, `wechatId`, `mobile`, `nickname`, `avatar`,
`gender`, `region`, `createTime`, `updateTime`
)
SELECT
t.wechatId AS identifier,
t.wechatId,
(SELECT phone FROM s2_wechat_friend WHERE wechatId = t.wechatId AND phone IS NOT NULL AND phone != '' LIMIT 1) AS mobile,
(SELECT nickname FROM s2_wechat_friend WHERE wechatId = t.wechatId ORDER BY id DESC LIMIT 1) AS nickname,
(SELECT avatar FROM s2_wechat_friend WHERE wechatId = t.wechatId ORDER BY id DESC LIMIT 1) AS avatar,
(SELECT gender FROM s2_wechat_friend WHERE wechatId = t.wechatId ORDER BY id DESC LIMIT 1) AS gender,
(SELECT region FROM s2_wechat_friend WHERE wechatId = t.wechatId ORDER BY id DESC LIMIT 1) AS region,
UNIX_TIMESTAMP() AS createTime,
UNIX_TIMESTAMP() AS updateTime
FROM (
SELECT wechatId FROM temp_wechat_ids LIMIT {$offset}, {$batchSize}
) AS t
ON DUPLICATE KEY UPDATE
mobile = COALESCE(VALUES(mobile), mobile),
nickname = COALESCE(VALUES(nickname), nickname),
avatar = COALESCE(VALUES(avatar), avatar),
gender = COALESCE(VALUES(gender), gender),
region = COALESCE(VALUES(region), region),
updateTime = UNIX_TIMESTAMP()";
$currentAffected = Db::execute($sql);
$affectedRows += $currentAffected;
if ($i % 5 == 0) {
gc_collect_cycles();
}
usleep(30000); // 30毫秒
}
// ========== 旧版流量池代码已废弃已迁移到V2 ==========
// for ($i = 0; $i < $batchCount; $i++) {
// $offset = $i * $batchSize;
//
// // 使用 ON DUPLICATE KEY UPDATE 支持插入和更新
// $sql = "INSERT INTO ck_traffic_pool_v1(
// `identifier`, `wechatId`, `mobile`, `nickname`, `avatar`,
// `gender`, `region`, `createTime`, `updateTime`
// )
// SELECT
// t.wechatId AS identifier,
// t.wechatId,
// (SELECT phone FROM s2_wechat_friend WHERE wechatId = t.wechatId AND phone IS NOT NULL AND phone != '' LIMIT 1) AS mobile,
// (SELECT nickname FROM s2_wechat_friend WHERE wechatId = t.wechatId ORDER BY id DESC LIMIT 1) AS nickname,
// (SELECT avatar FROM s2_wechat_friend WHERE wechatId = t.wechatId ORDER BY id DESC LIMIT 1) AS avatar,
// (SELECT gender FROM s2_wechat_friend WHERE wechatId = t.wechatId ORDER BY id DESC LIMIT 1) AS gender,
// (SELECT region FROM s2_wechat_friend WHERE wechatId = t.wechatId ORDER BY id DESC LIMIT 1) AS region,
// UNIX_TIMESTAMP() AS createTime,
// UNIX_TIMESTAMP() AS updateTime
// FROM (
// SELECT wechatId FROM temp_wechat_ids LIMIT {$offset}, {$batchSize}
// ) AS t
// ON DUPLICATE KEY UPDATE
// mobile = COALESCE(VALUES(mobile), mobile),
// nickname = COALESCE(VALUES(nickname), nickname),
// avatar = COALESCE(VALUES(avatar), avatar),
// gender = COALESCE(VALUES(gender), gender),
// region = COALESCE(VALUES(region), region),
// updateTime = UNIX_TIMESTAMP()";
//
// $currentAffected = Db::execute($sql);
// $affectedRows += $currentAffected;
//
// if ($i % 5 == 0) {
// gc_collect_cycles();
// }
//
// usleep(30000); // 30毫秒
// }
// ========== 旧版流量池代码结束(已迁移到 syncToTrafficPoolV2 ==========
// 注意:现在使用 syncToTrafficPoolV2() 方法同步到 V2 流量池系统
} catch (\Exception $e) {
\think\facade\Log::error("Error in traffic pool sync: " . $e->getMessage());
throw $e;