Bläddra i källkod

Вставка событий очередью

avolver 8 år sedan
förälder
incheckning
1b07adfe26

+ 12 - 0
src/Application.php

@@ -28,6 +28,14 @@ class Application extends ConsoleApplication
      */
     private $clickhouseClient;
 
+    /**
+     * {@inheritdoc}
+     */
+    public function __construct(string $name = 'Company events kitchen', string $version = '0.3.0')
+    {
+        parent::__construct($name, $version);
+    }
+
     /**
      * {@inheritdoc}
      */
@@ -40,6 +48,8 @@ class Application extends ConsoleApplication
     }
 
     /**
+     * Получение настроенного клиента ClickHouse
+     *
      * @return ClickHouseClient
      */
     public function getClickhouseClient(): ClickHouseClient
@@ -52,6 +62,8 @@ class Application extends ConsoleApplication
     }
 
     /**
+     * Получение настроенного клиента MongoDB
+     *
      * @return MongoDBClient
      */
     public function getMongoClient(): MongoDBClient

+ 3 - 3
src/Command/EventsRandomizerCommand.php

@@ -37,13 +37,13 @@ class EventsRandomizerCommand extends Command
     protected function execute(InputInterface $input, OutputInterface $output)
     {
         $count = $input->getOption('count');
-        $today = (new \DateTime())->setTime(0, 0);
+        $month = (new \DateTime(date('Y-m-01')))->setTime(0, 0);
 
         $companyRepo = new CompanyRepository($this->getApplication()->getMongoClient());
         $eventsRepo  = new EventsRepository(
             $this->getApplication()->getClickhouseClient(),
-            $today,
-            (clone $today)->add(new \DateInterval('P1M'))
+            $month,
+            (clone $month)->add(new \DateInterval('P1M'))->modify('-1 day')
         );
 
         $eventsRepo->generateEventsForCompanies(

+ 4 - 2
src/Repository/CompanyRepository.php

@@ -5,7 +5,6 @@ namespace Avolver\CompanyEventsKitchen\Repository;
 
 use Avolver\CompanyEventsKitchen\Traits\OutputTrait;
 use Symfony\Component\Console\Helper\ProgressBar;
-use Symfony\Component\Console\Output\OutputInterface;
 use MongoDB\Driver\Exception\Exception as MongoException;
 use MongoDB\Driver\Query as MongoQuery;
 use MongoDB\Driver\Cursor as MongoCursor;
@@ -73,7 +72,9 @@ class CompanyRepository
                 try {
                     $this->mongo->executeBulkWrite($bulk);
                 } catch (MongoException $e) {
-                    $this->output->writeln('Error: ' . $e->getMessage());
+                    if ($this->output) {
+                        $this->output->writeln('Error: ' . $e->getMessage());
+                    }
                 }
                 gc_collect_cycles();
                 unset($bulk);
@@ -99,6 +100,7 @@ class CompanyRepository
     public function findAll(): MongoCursor
     {
         $query  = new MongoQuery([]);
+
         return $this->mongo->getManager()->executeQuery(
             $this->mongo->getCompaniesNs(),
             $query

+ 48 - 16
src/Repository/EventsRepository.php

@@ -49,6 +49,13 @@ class EventsRepository
      */
     private $tableName;
 
+    /**
+     * Очередь событий для вставки
+     *
+     * @var array
+     */
+    private $bulkQueue = [];
+
     /**
      * EventsRepository constructor.
      *
@@ -108,50 +115,76 @@ QUERY;
      *
      * @throws \Exception
      */
-    public function generateEventsForCompanies(MongoCursor $cursor, int $needCount = 256): void
+    public function generateEventsForCompanies(MongoCursor $cursor, int $needCount = 3200): void
     {
-        $this->createTableIfNeeded();
+        //$this->createTableIfNeeded();
         $eventTypeCount = \count(CompanyEvent::keys());
         $this->clickhouse->database($this->databaseName);
 
+        $threshold = 100000;
+        $eachCount = 0;
+
         foreach ($cursor as $company) {
             $eventCount = random_int($needCount - 10, $needCount + 10);
             $eventCount = $eventCount > 0 ? $eventCount : 10;
             for ($count = 1; $count <= $eventCount; $count++) {
                 $randomDate  = $this->getRandomDate($this->startDate, $this->endDate);
                 $randomEventKey = random_int(1, $eventTypeCount);
-                $this->saveEvent($company->_id, $randomDate, new CompanyEvent($randomEventKey));
+                $this->addEventToQueue($company->_id, $randomDate, new CompanyEvent($randomEventKey));
+
+                $eachCount++;
+                if ($eachCount % $threshold === 0) {
+                    $this->insertEventsFromQueue();
+                    gc_collect_cycles();
+                    // echo " + 100k\n";
+                }
             }
+
         }
+
+        $this->insertEventsFromQueue();
     }
 
     /**
-     * Сохранение события
+     * Добавление события в очередь на сохранение
      *
      * @param ObjectId     $companyId
      * @param \DateTime    $datetime
      * @param CompanyEvent $event
      */
-    public function saveEvent(ObjectId $companyId, \DateTime $datetime, CompanyEvent $event): void
+    public function addEventToQueue(ObjectId $companyId, \DateTime $datetime, CompanyEvent $event): void
     {
         $idChunks = $this->convertObjectIdToInts($companyId);
 
+        $this->bulkQueue[] = [
+            $datetime->format('Y-m-d'),
+            $datetime->getTimestamp(),
+            $event->getKey(),
+            $idChunks[0],
+            $idChunks[1],
+            $idChunks[2]
+        ];
+    }
+
+    /**
+     * Вставка всех событий из очереди
+     */
+    public function insertEventsFromQueue(): void
+    {
+        if (!$this->bulkQueue || \count($this->bulkQueue)) {
+            return;
+        }
+
         $this->clickhouse->insert(
             $this->tableName,
             [
-                [
-                    $datetime->format('Y-m-d'),
-                    $datetime->getTimestamp(),
-                    $event->getKey(),
-                    $idChunks[0],
-                    $idChunks[1],
-                    $idChunks[2]
-                ]
+                $this->bulkQueue
             ],
             [
                 'event_date', 'event_time', 'event_type', 'id_1', 'id_2', 'id_3'
             ]
         );
+        $this->bulkQueue = [];
     }
 
     /**
@@ -185,12 +218,11 @@ QUERY;
     private function convertObjectIdToInts(ObjectId $objectId): array
     {
         $idString = (string) $objectId;
-        $idChunks = [
+
+        return [
             hexdec(mb_substr($idString, 0, 8)),
             hexdec(mb_substr($idString, 8, 8)),
             hexdec(mb_substr($idString, 16))
         ];
-
-        return $idChunks;
     }
 }