jobManager = $jobManager; $this->portionNumberProvider = $portionNumberProvider; $this->entityManager = $entityManager; } public function run(JobData $data): void { if (!$this->queue) { throw new RuntimeException("No queue name."); } $limit = $this->portionNumberProvider->get($this->queue); $group = $data->get('group'); $this->jobManager->processQueue($this->queue, $group, $limit); } public function prepare(ScheduledJobData $data, DateTimeImmutable $executeTime): void { $groupList = []; $shiftPeriod = self::SHIFT_PERIOD; $query = $this->entityManager ->getQueryBuilder() ->select('group') ->from(JobEntity::ENTITY_TYPE) ->where([ 'status' => JobStatus::PENDING, 'queue' => $this->queue, 'executeTime<=' => $executeTime ->modify($shiftPeriod) ->format(DateTime::SYSTEM_DATE_TIME_FORMAT), ]) ->groupBy('group') ->build(); $sth = $this->entityManager->getQueryExecutor()->execute($query); while ($row = $sth->fetch()) { $groupList[] = $row['group'] ?? null; } if (!count($groupList)) { return; } foreach ($groupList as $group) { $existingJob = $this->entityManager ->getRDBRepository(JobEntity::ENTITY_TYPE) ->select('id') ->where([ 'scheduledJobId' => $data->getId(), 'targetGroup' => $group, 'status' => [ JobStatus::RUNNING, JobStatus::READY, ] ]) ->findOne(); if ($existingJob) { continue; } $name = $data->getName(); if ($group) { $name .= ' :: ' . $group; } $this->entityManager->createEntity(JobEntity::ENTITY_TYPE, [ 'scheduledJobId' => $data->getId(), 'executeTime' => $executeTime->format(DateTime::SYSTEM_DATE_TIME_FORMAT), 'name' => $data->getName(), 'data' => [ 'group' => $group, ], 'targetGroup' => $group, ]); } } }