diff --git a/module/cron/control.php b/module/cron/control.php index 5877d5d108..b52506336a 100644 --- a/module/cron/control.php +++ b/module/cron/control.php @@ -137,35 +137,29 @@ class cron extends control { if(empty($this->config->global->cron)) return; - if('cli' !== PHP_SAPI) - { - ignore_user_abort(true); - set_time_limit(0); - session_write_close(); - } + ignore_user_abort(true); + set_time_limit(0); + session_write_close(); + + $execId = mt_rand(); + if($restart) $this->cron->restartCron($execId); $this->loadModel('common'); - $execId = mt_rand(); while(true) { - dao::$cache = array(); - - if($restart) $this->cron->clearQueue(); - - if($restart || $this->canSchedule($execId)) + $roles = $this->applyExecRoles($execId); + if(empty($roles)) { - $this->schedule($execId); - $this->execTasks($execId); + $this->cron->logCron('exit'); + ignore_user_abort(false); + return; + } - $restart = false; - sleep(20); - } - else - { - $this->execTasks($execId); - return; // If no task in queue, executor will exit. - } + if(in_array('scheduler', $roles)) $this->schedule($execId); + if(in_array('consumer', $roles)) $this->consumeTasks($execId); + + sleep(20); } } @@ -188,23 +182,17 @@ class cron extends control $execId = mt_rand(); while(true) { - dao::$cache = array(); - if(empty($this->config->global->cron)) { sleep(60); continue; } - if($restart) $this->cron->clearQueue(); + if($restart) $this->cron->restartCron($execId); - if($restart || $this->canSchedule($execId)) - { - $this->schedule($execId); - $restart = false; - } + if($this->canSchedule($execId)) $this->schedule($execId); - sleep(30); + sleep(20); } } @@ -223,6 +211,7 @@ class cron extends control $this->loadModel('common'); + $execId = mt_rand(); while(true) { if(empty($this->config->global->cron)) @@ -231,8 +220,8 @@ class cron extends control continue; } - $this->execTasks(mt_rand()); - sleep(10); + $this->consumeTasks($execId); + sleep(20); } } @@ -246,13 +235,57 @@ class cron extends control */ protected function canSchedule($execId) { - $settings = $this->dao->select('`key`,`value`')->from(TABLE_CONFIG)->where('owner')->eq('system')->andWhere('module')->eq('cron')->andWhere('section')->eq('run')->fetchPairs(); + $settings = $this->dao->select('`key`,`value`')->from(TABLE_CONFIG)->where('owner')->eq('system')->andWhere('module')->eq('cron')->andWhere('section')->eq('scheduler')->fetchPairs(); if(!isset($settings['execId']) || $settings['execId'] == $execId) return true; if(!isset($settings['lastTime']) || $settings['lastTime'] < date('Y-m-d H:i:s', strtotime('-1 minute'))) return true; return false; } + /** + * 检查该execId是否可以执行任务(最多允许1个调度进程,4个执行进程). + * Check the execId can exec(1 scheduler, 4 consumer). + * + * @param int $execId + * @access protected + * @return array + */ + protected function applyExecRoles($execId) + { + $roles = array(); + + $settings = $this->dao->select('*')->from(TABLE_CONFIG)->where('owner')->eq('system')->andWhere('module')->eq('cron')->fetchAll(); + + $scheduler = array('execId' => 0, 'lastTime' => ''); + $consumerCount = 0; + + $expirDate = date('Y-m-d H:i:s', strtotime('-1 minute')); + foreach($settings as $setting) + { + if($setting->section == 'scheduler' && $setting->key == 'lastTime') $scheduler['lastTime'] = $setting->value; + if($setting->section == 'scheduler' && $setting->key == 'execId') $scheduler['execId'] = $setting->value; + + if($setting->section == 'consumer') + { + if($setting->value > $expirDate) + { + if($setting->key == strval($execId)) $roles[] = 'consumer'; + $consumerCount ++; + } + else + { + $this->dao->delete()->from(TABLE_CONFIG)->where('id')->eq($setting->id)->exec(); + } + } + + } + + if(in_array($scheduler['execId'], array(0, $execId)) || $scheduler['lastTime'] < $expirDate) $roles[] = 'scheduler'; + if($consumerCount < 4) $roles[] = 'consumer'; + + return $roles; + } + /** * 调度生成队列任务. * Schedule, push tasks to queue. @@ -263,10 +296,9 @@ class cron extends control */ public function schedule($execId) { - $now = date(DT_DATETIME1); + dao::$cache = array(); - $this->loadModel('setting')->setItem('system.cron.run.execId', $execId); - $this->setting->setItem('system.cron.run.lastTime', $now); + $this->cron->updateTime('scheduler', $execId); /* Get and parse crons. */ $tasks = $this->dao->select('cron,MAX(`createdDate`) `datetime`')->from(TABLE_QUEUE)->groupBy('cron')->fetchAll('cron'); @@ -311,16 +343,19 @@ class cron extends control * @access public * @return bool */ - public function execTasks($execId) + public function consumeTasks($execId) { while(true) { dao::$cache = array(); + $this->cron->updateTime('consumer', $execId); + + /* Consume. */ $task = $this->dao->select('*')->from(TABLE_QUEUE)->where('status')->eq('wait')->andWhere('command')->ne('')->orderBy('createdDate')->fetch(); if(!$task) break; - $this->execTask($execId, $task); + $this->consumeTask($execId, $task); } } @@ -333,7 +368,7 @@ class cron extends control * @access public * @return bool */ - public function execTask($execId, $task) + public function consumeTask($execId, $task) { /* Other executor may execute the task at the same time,so we mark execId and wait 500ms to check whether we own it. */ $this->dao->update(TABLE_QUEUE)->set('status')->eq('doing')->set('execId')->eq($execId)->where('id')->eq($task->id)->exec(); diff --git a/module/cron/model.php b/module/cron/model.php index 3a4fe22798..8e4f248003 100644 --- a/module/cron/model.php +++ b/module/cron/model.php @@ -312,16 +312,91 @@ class cronModel extends model } /** - * Clear queue. + * Restart cron. * * @access public * @return void */ - public function clearQueue() + public function restartCron($execId) { + $this->dao->update(TABLE_CONFIG)->set('value')->eq($execId) + ->where('owner')->eq('system') + ->andWhere('module')->eq('cron') + ->andWhere('section')->eq('scheduler') + ->andWhere('`key`')->eq($execId) + ->exec(); $this->dao->delete()->from(TABLE_QUEUE)->where('createdDate')->le(date("Y-m-d H:i:s", strtotime("-1 week")))->exec(); - $log = date('G:i:s') . " clear\n\n"; - $this->logCron($log); + $this->logCron(date('G:i:s') . " restart\n\n"); + } + + /** + * Update last time of cron. + * + * @param string $role + * @param int $execId + * @access public + * @return void + */ + public function updateTime($role, $execId) + { + $now = date(DT_DATETIME1); + + $settings = $this->dao->select('*')->from(TABLE_CONFIG)->where('owner')->eq('system')->andWhere('module')->eq('cron')->andWhere('section')->eq($role)->fetchAll('key'); + if($role == 'scheduler') + { + if(isset($settings['execId'])) + { + $setting = $settings['execId']; + if($setting->value != strval($execId)) $this->dao->update(TABLE_CONFIG)->set('value')->eq($execId)->where('id')->eq($setting->id)->exec(); + } + else + { + $data = new stdclass(); + $data->owner = 'system'; + $data->module = 'cron'; + $data->section = 'scheduler'; + $data->key = 'execId'; + $data->value = $execId; + + $this->dao->insert(TABLE_CONFIG)->data($data)->exec(); + } + + if(isset($settings['lastTime'])) + { + $setting = $settings['lastTime']; + $this->dao->update(TABLE_CONFIG)->set('value')->eq($now)->where('id')->eq($setting->id)->exec(); + } + else + { + $data = new stdclass(); + $data->owner = 'system'; + $data->module = 'cron'; + $data->section = 'scheduler'; + $data->key = 'lastTime'; + $data->value = $now; + + $this->dao->insert(TABLE_CONFIG)->data($data)->exec(); + } + } + else + { + if(isset($settings[strval($execId)])) + { + $setting = $settings[strval($execId)]; + $this->dao->update(TABLE_CONFIG)->set('value')->eq($now)->where('id')->eq($setting->id)->exec(); + } + else + { + $data = new stdclass(); + $data->owner = 'system'; + $data->module = 'cron'; + $data->section = 'consumer'; + $data->key = $execId; + $data->value = $now; + + $this->dao->insert(TABLE_CONFIG)->data($data)->exec(); + } + } } }