Skip to content

Commit

Permalink
Deploy preview for PR 129 🛫
Browse files Browse the repository at this point in the history
  • Loading branch information
reidmeyer committed Aug 10, 2023
1 parent 7ddc8e0 commit d9997b3
Show file tree
Hide file tree
Showing 2 changed files with 6 additions and 14 deletions.
Binary file modified pr-preview/pr-129/sitemap.xml.gz
Binary file not shown.
20 changes: 6 additions & 14 deletions pr-preview/pr-129/stream/index.html
Original file line number Diff line number Diff line change
Expand Up @@ -1778,9 +1778,9 @@ <h2 id="kstreams.MetricsRebalanceListener" class="doc doc-heading">
<span class="sd"> to the consumer on the last rebalance</span>
<span class="sd"> &quot;&quot;&quot;</span>
<span class="c1"># lock all asyncio Tasks so no new metrics will be added to the Monitor</span>
<span class="k">if</span> <span class="n">revoked</span> <span class="ow">and</span> <span class="bp">self</span><span class="o">.</span><span class="n">engine</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
<span class="k">async</span> <span class="k">with</span> <span class="n">asyncio</span><span class="o">.</span><span class="n">Lock</span><span class="p">():</span>
<span class="k">await</span> <span class="bp">self</span><span class="o">.</span><span class="n">engine</span><span class="o">.</span><span class="n">monitor</span><span class="o">.</span><span class="n">stop</span><span class="p">()</span>
<span class="c1"># if revoked and self.engine is not None:</span>
<span class="c1"># async with asyncio.Lock():</span>
<span class="c1"># await self.engine.monitor.stop()</span>

<span class="k">async</span> <span class="k">def</span> <span class="nf">on_partitions_assigned</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">assigned</span><span class="p">:</span> <span class="n">Set</span><span class="p">[</span><span class="n">TopicPartition</span><span class="p">])</span> <span class="o">-&gt;</span> <span class="kc">None</span><span class="p">:</span>
<span class="w"> </span><span class="sd">&quot;&quot;&quot;</span>
Expand All @@ -1796,7 +1796,7 @@ <h2 id="kstreams.MetricsRebalanceListener" class="doc doc-heading">
<span class="c1"># lock all asyncio Tasks so no new metrics will be added to the Monitor</span>
<span class="k">if</span> <span class="n">assigned</span> <span class="ow">and</span> <span class="bp">self</span><span class="o">.</span><span class="n">engine</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
<span class="k">async</span> <span class="k">with</span> <span class="n">asyncio</span><span class="o">.</span><span class="n">Lock</span><span class="p">():</span>
<span class="bp">self</span><span class="o">.</span><span class="n">engine</span><span class="o">.</span><span class="n">monitor</span><span class="o">.</span><span class="n">start</span><span class="p">()</span>
<span class="c1"># self.engine.monitor.start()</span>

<span class="n">stream</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">stream</span>
<span class="k">if</span> <span class="n">stream</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
Expand Down Expand Up @@ -1898,7 +1898,7 @@ <h3 id="kstreams.rebalance_listener.MetricsRebalanceListener.on_partitions_assig
<span class="c1"># lock all asyncio Tasks so no new metrics will be added to the Monitor</span>
<span class="k">if</span> <span class="n">assigned</span> <span class="ow">and</span> <span class="bp">self</span><span class="o">.</span><span class="n">engine</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
<span class="k">async</span> <span class="k">with</span> <span class="n">asyncio</span><span class="o">.</span><span class="n">Lock</span><span class="p">():</span>
<span class="bp">self</span><span class="o">.</span><span class="n">engine</span><span class="o">.</span><span class="n">monitor</span><span class="o">.</span><span class="n">start</span><span class="p">()</span>
<span class="c1"># self.engine.monitor.start()</span>

<span class="n">stream</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">stream</span>
<span class="k">if</span> <span class="n">stream</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
Expand Down Expand Up @@ -1969,11 +1969,7 @@ <h3 id="kstreams.rebalance_listener.MetricsRebalanceListener.on_partitions_revok
<span class="normal">108</span>
<span class="normal">109</span>
<span class="normal">110</span>
<span class="normal">111</span>
<span class="normal">112</span>
<span class="normal">113</span>
<span class="normal">114</span>
<span class="normal">115</span></pre></div></td><td class="code"><div><pre><span></span><code><span class="k">async</span> <span class="k">def</span> <span class="nf">on_partitions_revoked</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">revoked</span><span class="p">:</span> <span class="n">Set</span><span class="p">[</span><span class="n">TopicPartition</span><span class="p">])</span> <span class="o">-&gt;</span> <span class="kc">None</span><span class="p">:</span>
<span class="normal">111</span></pre></div></td><td class="code"><div><pre><span></span><code><span class="k">async</span> <span class="k">def</span> <span class="nf">on_partitions_revoked</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">revoked</span><span class="p">:</span> <span class="n">Set</span><span class="p">[</span><span class="n">TopicPartition</span><span class="p">])</span> <span class="o">-&gt;</span> <span class="kc">None</span><span class="p">:</span>
<span class="w"> </span><span class="sd">&quot;&quot;&quot;</span>
<span class="sd"> Coroutine to be called *before* a rebalance operation starts and</span>
<span class="sd"> *after* the consumer stops fetching data.</span>
Expand All @@ -1984,10 +1980,6 @@ <h3 id="kstreams.rebalance_listener.MetricsRebalanceListener.on_partitions_revok
<span class="sd"> revoked Set[TopicPartitions]: Partitions that were assigned</span>
<span class="sd"> to the consumer on the last rebalance</span>
<span class="sd"> &quot;&quot;&quot;</span>
<span class="c1"># lock all asyncio Tasks so no new metrics will be added to the Monitor</span>
<span class="k">if</span> <span class="n">revoked</span> <span class="ow">and</span> <span class="bp">self</span><span class="o">.</span><span class="n">engine</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
<span class="k">async</span> <span class="k">with</span> <span class="n">asyncio</span><span class="o">.</span><span class="n">Lock</span><span class="p">():</span>
<span class="k">await</span> <span class="bp">self</span><span class="o">.</span><span class="n">engine</span><span class="o">.</span><span class="n">monitor</span><span class="o">.</span><span class="n">stop</span><span class="p">()</span>
</code></pre></div></td></tr></table></div>
</details>
</div>
Expand Down

0 comments on commit d9997b3

Please sign in to comment.