1 | <?php |
||
32 | final class FirehoseAdapter implements AdapterInterface |
||
33 | { |
||
34 | const BATCHSIZE_SEND = 100; |
||
35 | |||
36 | /** @var FirehoseClient */ |
||
37 | protected $client; |
||
38 | |||
39 | /** @var array */ |
||
40 | protected $options; |
||
41 | |||
42 | /** @var string */ |
||
43 | protected $deliveryStreamName; |
||
44 | |||
45 | /** |
||
46 | * @param FirehoseClient $client |
||
47 | * @param string $deliveryStreamName |
||
48 | * @param array $options - BatchSize <integer> The number of messages to send in each batch. |
||
49 | */ |
||
50 | 8 | public function __construct(FirehoseClient $client, $deliveryStreamName, array $options = []) |
|
56 | |||
57 | /** |
||
58 | * @param MessageInterface[] $messages |
||
59 | * |
||
60 | * @throws MethodNotSupportedException |
||
61 | */ |
||
62 | 1 | public function acknowledge(array $messages) |
|
70 | |||
71 | /** |
||
72 | * @param MessageInterface[] $messages |
||
73 | */ |
||
74 | public function reject(array $messages) |
||
82 | |||
83 | /** |
||
84 | * @param MessageFactoryInterface $factory |
||
85 | * @param int $limit |
||
86 | * |
||
87 | * @throws MethodNotSupportedException |
||
88 | */ |
||
89 | 1 | public function dequeue(MessageFactoryInterface $factory, $limit) |
|
90 | { |
||
91 | 1 | throw new MethodNotSupportedException( |
|
92 | 1 | __FUNCTION__, |
|
93 | 1 | $this, |
|
94 | 1 | [] |
|
95 | ); |
||
96 | } |
||
97 | |||
98 | /** |
||
99 | * @param MessageInterface[] $messages |
||
100 | * |
||
101 | * @throws FailedEnqueueException |
||
102 | */ |
||
103 | 3 | public function enqueue(array $messages) |
|
104 | { |
||
105 | 3 | $failed = []; |
|
106 | 3 | $batches = array_chunk( |
|
107 | 3 | $messages, |
|
108 | 3 | $this->getOption('BatchSize', self::BATCHSIZE_SEND) |
|
109 | ); |
||
110 | |||
111 | 3 | foreach ($batches as $batch) { |
|
112 | 3 | $requestRecords = array_map( |
|
113 | 3 | function (MessageInterface $message) { |
|
114 | return [ |
||
115 | 3 | 'Data' => $message->getBody(), |
|
116 | ]; |
||
117 | 3 | }, |
|
118 | 3 | $batch |
|
119 | ); |
||
120 | |||
121 | $request = [ |
||
122 | 3 | 'DeliveryStreamName' => $this->deliveryStreamName, |
|
123 | 3 | 'Records' => $requestRecords, |
|
124 | ]; |
||
125 | |||
126 | 3 | $results = $this->client->putRecordBatch($request); |
|
127 | |||
128 | 3 | foreach ($results->get('RequestResponses') as $idx => $response) { |
|
129 | 1 | if (isset($response['ErrorCode'])) { |
|
130 | 3 | $failed[] = $batch[$idx]; |
|
131 | } |
||
132 | } |
||
133 | } |
||
134 | |||
135 | 3 | if (!empty($failed)) { |
|
136 | 1 | throw new FailedEnqueueException($this, $failed); |
|
137 | } |
||
138 | 2 | } |
|
139 | |||
140 | /** |
||
141 | * @param string $name |
||
142 | * @param mixed $default |
||
143 | * |
||
144 | * @return mixed |
||
145 | */ |
||
146 | 3 | protected function getOption($name, $default = null) |
|
150 | |||
151 | /** |
||
152 | * @throws MethodNotSupportedException |
||
153 | */ |
||
154 | 1 | public function purge() |
|
162 | |||
163 | /** |
||
164 | * @throws MethodNotSupportedException |
||
165 | */ |
||
166 | 1 | public function delete() |
|
174 | } |
||
175 |