attempts());
//通过这个方法可以检查这个任务已经重试了几次了
if ($job->attempts() > 5) {
Log::info("# 重试次数已达上限 {$job->attempts()}");
$job->delete();
return;
}
// 如有必要,可以根据业务需求和数据库中的最新数据,判断该任务是否仍有必要执行.
$isJobStillNeedToBeDone = $this->checkDatabaseToSeeIfJobNeedToBeDone($data);
if ($isJobStillNeedToBeDone) {
Log::info("# 发现通知已完成,队列 deleted");
$job->delete();
return;
}
$isJobDone = $this->doJob($data);
if ($isJobDone) {
//如果任务执行成功, 记得删除任务
$job->delete();
Log::info("# 订单处理完成,队列 deleted");
//print("Hello Job has been done and deleted" . "\n");
} else {
if ($job->attempts() > 5) {
//通过这个方法可以检查这个任务已经重试了几次了
Log::info("# 重试次数已达上限,当前 {$job->attempts()} 次");
//print("Hello Job has been retried more than 5 times!" . "\n");
$job->delete();
// 也可以重新发布这个任务
//print("Hello Job will be availabe again after 2s."."\n");
//$job->release(2); //$delay为延迟时间,表示该任务延迟2秒后再执行
} else {
Log::info("# 订单重试,当前 {$job->attempts()} 次");
$job->release(120); //$delay为延迟时间,表示该任务延迟120秒后再执行
}
}
}
/**
* 有些消息在到达消费者时,可能已经不再需要执行了
* @param array|mixed $data 发布任务时自定义的数据
* @return boolean 任务执行的结果
* @throws \think\db\exception\DataNotFoundException
* @throws \think\db\exception\ModelNotFoundException
* @throws \think\exception\DbException
*/
private function checkDatabaseToSeeIfJobNeedToBeDone($data)
{
$payCpinfoModel = new PayCpinfo;
//1 查询数据库判断是否已经更新
$where['orderid'] = $data['orderid'];
$where['status'] = 0;
$where['payflag'] = 1;
$this->data = $info = $payCpinfoModel->where($where)->find();
if ($info) {
return false;
}
return true;
}
/**
* 根据消息中的数据进行实际的业务处理
* @param array|mixed $data 发布任务时自定义的数据
* @return boolean 任务执行的结果
*/
private function doJob($data)
{
// 将字符串参数转为数组形式
if (!is_array($data['params'])) parse_str($data['params'], $data['params']);
$guzzle = new Client();
try {
$result = $guzzle->request($data['http_method'], $data['url'], [
'form_params' => $data['params'],
'timeout' => 30,
]);
} catch (GuzzleException $exception) {
Log::info("# doJob {$data['orderid']} 订单通知失败: " . $exception->getMessage());
//todo 错误信息记录
return false;
}
if (0 != strcasecmp($result->getBody(), $data['success_flag'])) {
return false;
}
return $this->update();
}
private function update()
{
$payCpinfoModel = new PayCpinfo;
$res = $payCpinfoModel->where(['id' => $this->data['id']])->update([
'status' => 1,
'update_time' => time(),
]);
if (!$res) {
return false;
}
return true;
}
}