blob: beb68cf2e257af9c97c5ad98f5aeacf244bc2dbf [file]
<?xml version="1.0" encoding="utf-8"?>
<!DOCTYPE html PUBLIC "-//W3C//DTD XHTML 1.0 Strict//EN"
"DTD/xhtml1-strict.dtd">
<html>
<head>
<title>pulsar.asyncio.Client</title>
<meta name="generator" content="pydoctor 25.4.0">
</meta>
<meta http-equiv="Content-Type" content="text/html;charset=utf-8" />
<meta name="viewport" content="width=device-width, initial-scale=1 maximum-scale=1" />
<link rel="stylesheet" type="text/css" href="apidocs.css" />
<link rel="stylesheet" type="text/css" href="readthedocstheme.css" />
<link rel="stylesheet" type="text/css" href="extra.css" />
</head>
<body>
<nav class="navbar navbar-default mainnavbar">
<div class="container-fluid">
<div class="navbar-header">
<div class="navlinks">
<span class="navbar-brand">
pulsar/_pulsar <a href="index.html">API Documentation</a>
</span>
<a href="moduleIndex.html">
Modules
</a>
<a href="classIndex.html">
Classes
</a>
<a href="nameIndex.html">
Names
</a>
<div id="search-box-container">
<div class="input-group">
<input id="search-box" type="search" name="search-query" placeholder="Search..." aria-label="Search" minlength="2" class="form-control" autocomplete="off" />
<span class="input-group-btn">
<a style="display: none;" class="btn btn-default" id="search-clear-button" title="Clear" onclick="clearSearch()"><img src="fonts/x-circle.svg" alt="Clear" /></a>
<a class="btn btn-default" id="apidocs-help-button" title="Help" href="apidocs-help.html"><img src="fonts/info.svg" alt="Help" /></a>
</span>
</div>
</div>
</div>
<div id="search-results-container" style="display: none;">
<div id="search-buttons">
<span class="label label-default" id="search-docstrings-button">
<label class="checkbox-inline">
<input type="checkbox" id="toggle-search-in-docstrings-checkbox" value="false" onclick="toggleSearchInDocstrings()">
search in docstrings
</input>
</label>
</span>
</div>
<noscript>
<h1>Cannot search: JavaScript is not supported/enabled in your browser.</h1>
</noscript>
<div id="search-status"> </div>
<div class="warning" id="search-warn-box" style="display: none;">
<p class="rst-last"><span id="search-warn"></span></p>
</div>
<table id="search-results">
<!-- Filled dynamically by JS -->
</table>
<div style="margin-top: 10px;">
<p>For more information on the search, visit the <a href="apidocs-help.html#rst-search">help page</a>.</p>
</div>
</div>
</div>
</div>
<!-- Side navigation -->
<div class="sidebarcontainer">
<nav class="sidebar">
<div>
<div class="thingTitle">
<span>Class</span>
<code class="thisobject"><a href="pulsar.asyncio.Client.html" class="internal-link" title="This class"><wbr></wbr>Client</a></code>
</div>
<div>
<div class="childrenKindTitle">Methods</div>
<ul>
<li class="">
<div class="itemName"><code><a href="#__init__" class="internal-link" title="pulsar.asyncio.Client.__init__">__init__</a></code>
</div>
</li><li class="">
<div class="itemName"><code><a href="#close" class="internal-link" title="pulsar.asyncio.Client.close">close</a></code>
</div>
</li><li class="">
<div class="itemName"><code><a href="#create_producer" class="internal-link" title="pulsar.asyncio.Client.create_producer">create<wbr></wbr>_producer</a></code>
</div>
</li><li class="">
<div class="itemName"><code><a href="#subscribe" class="internal-link" title="pulsar.asyncio.Client.subscribe">subscribe</a></code>
</div>
</li>
</ul>
<div class="childrenKindTitle">Attributes</div>
<ul>
<li class="private">
<div class="itemName"><code><a href="#_client" class="internal-link" title="pulsar.asyncio.Client._client">_client</a></code>
</div>
</li>
</ul>
</div>
</div><div>
<div class="thingTitle">
<span>Module</span>
<code><a href="pulsar.asyncio.html" class="internal-link" title="The parent of this class">asyncio</a></code>
</div>
<div>
<div class="childrenKindTitle">Classes</div>
<ul>
<li class=" thisobject">
<div class="itemName"><code><a href="pulsar.asyncio.Client.html" class="internal-link" title="pulsar.asyncio.Client"><wbr></wbr>Client</a></code>
</div>
</li><li class="">
<div class="itemName"><code><a href="pulsar.asyncio.Consumer.html" class="internal-link" title="pulsar.asyncio.Consumer"><wbr></wbr>Consumer</a></code>
</div>
</li><li class="">
<div class="itemName"><code><a href="pulsar.asyncio.Producer.html" class="internal-link" title="pulsar.asyncio.Producer"><wbr></wbr>Producer</a></code>
</div>
</li><li class="">
<div class="itemName"><code><a href="pulsar.asyncio.PulsarException.html" class="internal-link" title="pulsar.asyncio.PulsarException"><wbr></wbr>Pulsar<wbr></wbr>Exception</a></code>
</div>
</li>
</ul>
<div class="childrenKindTitle">Functions</div>
<ul>
<li class="private">
<div class="itemName"><code><a href="pulsar.asyncio.html#_set_future" class="internal-link" title="pulsar.asyncio._set_future">_set<wbr></wbr>_future</a></code>
</div>
</li>
</ul>
</div>
</div>
</nav>
<!-- No sidebar toggle for read the docs theme, the sidebar is always
visible when the screen is width enough.
-->
</div>
</nav>
<div class="container-fluid">
<div id="main" class="">
<div class="page-header">
<h1 class="class"><code><code><a href="pulsar.html" class="internal-link">pulsar</a></code><wbr></wbr>.<code><a href="pulsar.asyncio.html" class="internal-link" title="pulsar.asyncio">asyncio</a></code><wbr></wbr>.<code><a href="pulsar.asyncio.Client.html" class="internal-link" title="pulsar.asyncio.Client">Client</a></code></code></h1>
<div id="showPrivate">
<button class="btn btn-link" onclick="togglePrivate()">Toggle Private API</button>
</div>
</div>
<div class="categoryHeader">
class documentation
</div>
<div class="extrasDocstring">
<p class="class-signature"><code><span class="py-keyword">class</span> <span class="py-defname">Client</span>: <a href="https://github.com/apache/pulsar-client-python/tree/v3.9.0/pulsar/asyncio.py#L395" class="sourceLink">(source)</a></code></p><p>Constructor: <code><a href="#__init__" class="internal-link" title="pulsar.asyncio.Client.__init__">Client(service_url, **kwargs)</a></code></p>
<p><a href="classIndex.html#pulsar.asyncio.Client">View In Hierarchy</a></p>
</div>
<div class="moduleDocstring">
<div><p>The asynchronous version of <code><a href="pulsar.Client.html" class="internal-link">pulsar.Client</a></code>.</p></div>
</div>
<div id="splitTables">
<table class="children sortable" id="id9">
<tr class="method">
<td>Method</td>
<td><code><a href="#__init__" class="internal-link" title="pulsar.asyncio.Client.__init__">__init__</a></code></td>
<td>See <code><a href="pulsar.Client.html#__init__" class="internal-link">pulsar.Client.__init__</a></code></td>
</tr><tr class="method">
<td>Async Method</td>
<td><code><a href="#close" class="internal-link" title="pulsar.asyncio.Client.close">close</a></code></td>
<td>Close the client and all the associated producers and consumers</td>
</tr><tr class="method">
<td>Async Method</td>
<td><code><a href="#create_producer" class="internal-link" title="pulsar.asyncio.Client.create_producer">create<wbr></wbr>_producer</a></code></td>
<td>Create a new producer on a given topic</td>
</tr><tr class="method">
<td>Async Method</td>
<td><code><a href="#subscribe" class="internal-link" title="pulsar.asyncio.Client.subscribe">subscribe</a></code></td>
<td>Subscribe to the given topic and subscription combination.</td>
</tr><tr class="instancevariable private">
<td>Instance Variable</td>
<td><code><a href="#_client" class="internal-link" title="pulsar.asyncio.Client._client">_client</a></code></td>
<td><span class="rst-undocumented">Undocumented</span></td>
</tr>
</table>
</div>
<div id="childList">
<div class="basemethod">
<a name="pulsar.asyncio.Client.__init__">
</a>
<a name="__init__">
</a>
<div class="functionHeader">
<span class="py-keyword">def</span>&#160;<span class="py-defname">__init__</span><span class="function-signature">(<span class="rst-sig-param"><span class="rst-undocumented">self</span>, </span><span class="rst-sig-param">service_url, </span><span class="rst-sig-param">**kwargs</span>):</span>
<a class="sourceLink" href="https://github.com/apache/pulsar-client-python/tree/v3.9.0/pulsar/asyncio.py#L400">
(source)
</a>
<a class="headerLink" href="#__init__" title="pulsar.asyncio.Client.__init__">
ΒΆ
</a>
</div>
<div class="docstring functionBody">
<div><p>See <code><a href="pulsar.Client.html#__init__" class="internal-link">pulsar.Client.__init__</a></code></p></div>
</div>
</div><div class="basemethod">
<a name="pulsar.asyncio.Client.close">
</a>
<a name="close">
</a>
<div class="functionHeader">
<span class="py-keyword">async&#160;def</span>&#160;<span class="py-defname">close</span><span class="function-signature">(<span class="rst-sig-param"><span class="rst-undocumented">self</span></span>):</span>
<a class="sourceLink" href="https://github.com/apache/pulsar-client-python/tree/v3.9.0/pulsar/asyncio.py#L728">
(source)
</a>
<a class="headerLink" href="#close" title="pulsar.asyncio.Client.close">
ΒΆ
</a>
</div>
<div class="docstring functionBody">
<div><p>Close the client and all the associated producers and consumers</p><table class="fieldTable"><tr class="fieldStart"><td class="fieldName" colspan="2">Raises</td></tr><tr><td class="fieldArgContainer"><code><a href="pulsar.asyncio.PulsarException.html" class="internal-link" title="pulsar.asyncio.PulsarException">PulsarException</a></code></td><td></td></tr></table></div>
</div>
</div><div class="basemethod">
<a name="pulsar.asyncio.Client.create_producer">
</a>
<a name="create_producer">
</a>
<div class="functionHeader">
<span class="py-keyword">async&#160;def</span>&#160;<span class="py-defname">create_producer</span><span class="function-signature long-signature">(<span class="rst-sig-param"><span class="rst-undocumented">self</span>, </span><span class="rst-sig-param">topic: <code><a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a></code>, </span><span class="rst-sig-param">producer_name: <code><a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">schema: <code><a href="pulsar.schema.schema.Schema.html" class="internal-link" title="pulsar.schema.schema.Schema">pulsar.schema.Schema</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">initial_sequence_id: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">send_timeout_millis: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 30000, </span><span class="rst-sig-param">compression_type: <code><a href="_pulsar.CompressionType.html" class="internal-link" title="_pulsar.CompressionType">CompressionType</a></code> = CompressionType.NONE, </span><span class="rst-sig-param">max_pending_messages: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 1000, </span><span class="rst-sig-param">max_pending_messages_across_partitions: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 50000, </span><span class="rst-sig-param">block_if_queue_full: <code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code> = <a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a>, </span><span class="rst-sig-param">batching_enabled: <code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code> = <a href="https://docs.python.org/3/library/constants.html#True" class="intersphinx-link">True</a>, </span><span class="rst-sig-param">batching_max_messages: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 1000, </span><span class="rst-sig-param">batching_max_allowed_size_in_bytes: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 128 * 1024, </span><span class="rst-sig-param">batching_max_publish_delay_ms: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 10, </span><span class="rst-sig-param">chunking_enabled: <code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code> = <a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a>, </span><span class="rst-sig-param">message_routing_mode: <code><a href="_pulsar.PartitionsRoutingMode.html" class="internal-link" title="_pulsar.PartitionsRoutingMode">PartitionsRoutingMode</a></code> = PartitionsRoutingMode.RoundRobinDistribution, </span><span class="rst-sig-param">lazy_start_partitioned_producers: <code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code> = <a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a>, </span><span class="rst-sig-param">properties: <code><a href="https://docs.python.org/3/library/stdtypes.html#dict" class="intersphinx-link">dict</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">batching_type: <code><a href="_pulsar.BatchingType.html" class="internal-link" title="_pulsar.BatchingType">BatchingType</a></code> = BatchingType.Default, </span><span class="rst-sig-param">encryption_key: <code><a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">crypto_key_reader: <code><a href="pulsar.CryptoKeyReader.html" class="internal-link" title="pulsar.CryptoKeyReader">pulsar.CryptoKeyReader</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">access_mode: <code><a href="_pulsar.ProducerAccessMode.html" class="internal-link" title="_pulsar.ProducerAccessMode">ProducerAccessMode</a></code> = ProducerAccessMode.Shared, </span><span class="rst-sig-param">message_router: <code><a href="https://docs.python.org/3/library/typing.html#typing.Callable" class="intersphinx-link">Callable</a>[<wbr></wbr>[<wbr></wbr><a href="pulsar.Message.html" class="internal-link" title="pulsar.Message">pulsar.Message</a>, <wbr></wbr><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a>], <wbr></wbr><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a>] | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></span>) -&gt; <code><a href="pulsar.asyncio.Producer.html" class="internal-link" title="pulsar.asyncio.Producer">Producer</a></code>:</span>
<a class="sourceLink" href="https://github.com/apache/pulsar-client-python/tree/v3.9.0/pulsar/asyncio.py#L407">
(source)
</a>
<a class="headerLink" href="#create_producer" title="pulsar.asyncio.Client.create_producer">
ΒΆ
</a>
</div>
<div class="docstring functionBody">
<div><p>Create a new producer on a given topic</p><table class="fieldTable"><tr class="fieldStart"><td class="fieldName" colspan="2">Parameters</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">topic:</span><code><a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a></code></td><td class="fieldArgDesc">The topic name</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">producer<wbr></wbr>_name:</span>str|<blockquote>
None</blockquote>
</td><td class="fieldArgDesc">Specify a name for the producer. If not assigned, the system will generate a globally
unique name which can be accessed with <code><a href="pulsar.asyncio.Producer.html#producer_name" class="internal-link" title="pulsar.asyncio.Producer.producer_name">Producer.producer_name()</a></code>. When specifying a
name, it is up to the user to ensure that, for a given topic, the producer name is
unique across all Pulsar's clusters.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">schema:</span>pulsar.schema.Schema|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Define the schema of the data that will be published by this producer.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">initial<wbr></wbr>_sequence<wbr></wbr>_id:</span>int|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Set the baseline for the sequence ids for messages published by
the producer.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">send<wbr></wbr>_timeout<wbr></wbr>_millis:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> <span class="literal">30000</span></td><td class="fieldArgDesc">If a message is not acknowledged by the server before the
send_timeout expires, an error will be reported.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">compression<wbr></wbr>_type:</span><code><a href="_pulsar.CompressionType.html" class="internal-link" title="_pulsar.CompressionType">CompressionType</a></code>, <em>default</em> <code>CompressionType.NONE</code></td><td class="fieldArgDesc">Set the compression type for the producer.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">max<wbr></wbr>_pending<wbr></wbr>_messages:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> <span class="literal">1000</span></td><td class="fieldArgDesc">Set the max size of the queue holding the messages pending to
receive an acknowledgment from the broker.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">max<wbr></wbr>_pending<wbr></wbr>_messages<wbr></wbr>_across<wbr></wbr>_partitions:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> <span class="literal">50000</span></td><td class="fieldArgDesc">Set the max size of the queue holding the messages pending to
receive an acknowledgment across partitions.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">block<wbr></wbr>_if<wbr></wbr>_queue<wbr></wbr>_full:</span><code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code>, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a></code></td><td class="fieldArgDesc">Set whether send operations should block when the outgoing
message queue is full.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">batching<wbr></wbr>_enabled:</span><code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code>, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#True" class="intersphinx-link">True</a></code></td><td class="fieldArgDesc">Enable automatic message batching. Note that, unlike the synchronous producer API in
<tt class="rst-docutils rst-literal">pulsar.Client.create_producer</tt>, batching is enabled by default for the asyncio
producer.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">batching<wbr></wbr>_max<wbr></wbr>_messages:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> <span class="literal">1000</span></td><td class="fieldArgDesc">Maximum number of messages in a batch.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">batching<wbr></wbr>_max<wbr></wbr>_allowed<wbr></wbr>_size<wbr></wbr>_in<wbr></wbr>_bytes:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> 128*1024</td><td class="fieldArgDesc">Maximum size in bytes of a batch.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">batching<wbr></wbr>_max<wbr></wbr>_publish<wbr></wbr>_delay<wbr></wbr>_ms:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> <span class="literal">10</span></td><td class="fieldArgDesc">The batch interval in milliseconds.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">chunking<wbr></wbr>_enabled:</span><code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code>, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a></code></td><td class="fieldArgDesc">Enable chunking of large messages.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">message<wbr></wbr>_routing<wbr></wbr>_mode:</span><code><a href="_pulsar.PartitionsRoutingMode.html" class="internal-link" title="_pulsar.PartitionsRoutingMode">PartitionsRoutingMode</a></code>, </td><td class="fieldArgDesc">default=PartitionsRoutingMode.RoundRobinDistribution
Set the message routing mode for the partitioned producer.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">lazy<wbr></wbr>_start<wbr></wbr>_partitioned<wbr></wbr>_producers:</span><code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code>, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a></code></td><td class="fieldArgDesc">Start partitioned producers lazily on demand.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">properties:</span>dict|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Sets the properties for the producer.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">batching<wbr></wbr>_type:</span><code><a href="_pulsar.BatchingType.html" class="internal-link" title="_pulsar.BatchingType">BatchingType</a></code>, <em>default</em> <code>BatchingType.Default</code></td><td class="fieldArgDesc">Sets the batching type for the producer.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">encryption<wbr></wbr>_key:</span>str|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">The key used for symmetric encryption.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">crypto<wbr></wbr>_key<wbr></wbr>_reader:</span>pulsar.CryptoKeyReader|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Symmetric encryption class implementation.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">access<wbr></wbr>_mode:</span><code><a href="_pulsar.ProducerAccessMode.html" class="internal-link" title="_pulsar.ProducerAccessMode">ProducerAccessMode</a></code>, <em>default</em> <code>ProducerAccessMode.Shared</code></td><td class="fieldArgDesc">Set the type of access mode that the producer requires on the topic.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">message<wbr></wbr>_router:</span><code><a href="https://docs.python.org/3/library/typing.html#typing.Callable" class="intersphinx-link">Callable</a>[[<a href="pulsar.Message.html" class="internal-link">pulsar.Message</a></code>, <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a>]</code>, <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a>]</code> |<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">A custom message router function that takes a Message and the
number of partitions and returns the partition index.</td></tr></table><table class="fieldTable"><tr class="fieldStart"><td class="fieldName" colspan="2">Returns</td></tr><tr><td class="fieldArgContainer"><code><a href="pulsar.asyncio.Producer.html" class="internal-link" title="pulsar.asyncio.Producer">Producer</a></code></td><td class="fieldArgDesc">The producer created</td></tr></table><table class="fieldTable"><tr class="fieldStart"><td class="fieldName" colspan="2">Raises</td></tr><tr><td class="fieldArgContainer"><code><a href="pulsar.asyncio.PulsarException.html" class="internal-link" title="pulsar.asyncio.PulsarException">PulsarException</a></code></td><td></td></tr></table></div>
</div>
</div><div class="basemethod">
<a name="pulsar.asyncio.Client.subscribe">
</a>
<a name="subscribe">
</a>
<div class="functionHeader">
<span class="py-keyword">async&#160;def</span>&#160;<span class="py-defname">subscribe</span><span class="function-signature long-signature">(<span class="rst-sig-param"><span class="rst-undocumented">self</span>, </span><span class="rst-sig-param">topic: <code><a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a> | <a href="https://docs.python.org/3/library/stdtypes.html#list" class="intersphinx-link">list</a>[<wbr></wbr><a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a>]</code>, </span><span class="rst-sig-param">subscription_name: <code><a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a></code>, </span><span class="rst-sig-param">consumer_type: <code><a href="_pulsar.ConsumerType.html" class="internal-link" title="_pulsar.ConsumerType">pulsar.ConsumerType</a></code> = pulsar.ConsumerType.Exclusive, </span><span class="rst-sig-param">schema: <code><a href="pulsar.schema.schema.Schema.html" class="internal-link" title="pulsar.schema.schema.Schema">pulsar.schema.Schema</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">receiver_queue_size: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 1000, </span><span class="rst-sig-param">max_total_receiver_queue_size_across_partitions: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 50000, </span><span class="rst-sig-param">consumer_name: <code><a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">unacked_messages_timeout_ms: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">broker_consumer_stats_cache_time_ms: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 30000, </span><span class="rst-sig-param">negative_ack_redelivery_delay_ms: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 60000, </span><span class="rst-sig-param">is_read_compacted: <code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code> = <a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a>, </span><span class="rst-sig-param">properties: <code><a href="https://docs.python.org/3/library/stdtypes.html#dict" class="intersphinx-link">dict</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">initial_position: <code><a href="_pulsar.InitialPosition.html" class="internal-link" title="_pulsar.InitialPosition">InitialPosition</a></code> = InitialPosition.Latest, </span><span class="rst-sig-param">crypto_key_reader: <code><a href="pulsar.CryptoKeyReader.html" class="internal-link" title="pulsar.CryptoKeyReader">pulsar.CryptoKeyReader</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">replicate_subscription_state_enabled: <code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code> = <a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a>, </span><span class="rst-sig-param">max_pending_chunked_message: <code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code> = 10, </span><span class="rst-sig-param">auto_ack_oldest_chunked_message_on_queue_full: <code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code> = <a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a>, </span><span class="rst-sig-param">start_message_id_inclusive: <code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code> = <a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a>, </span><span class="rst-sig-param">batch_receive_policy: <code><a href="pulsar.ConsumerBatchReceivePolicy.html" class="internal-link" title="pulsar.ConsumerBatchReceivePolicy">pulsar.ConsumerBatchReceivePolicy</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">key_shared_policy: <code><a href="pulsar.ConsumerKeySharedPolicy.html" class="internal-link" title="pulsar.ConsumerKeySharedPolicy">pulsar.ConsumerKeySharedPolicy</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">batch_index_ack_enabled: <code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code> = <a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a>, </span><span class="rst-sig-param">regex_subscription_mode: <code><a href="_pulsar.RegexSubscriptionMode.html" class="internal-link" title="_pulsar.RegexSubscriptionMode">RegexSubscriptionMode</a></code> = RegexSubscriptionMode.PersistentOnly, </span><span class="rst-sig-param">dead_letter_policy: <code><a href="pulsar.ConsumerDeadLetterPolicy.html" class="internal-link" title="pulsar.ConsumerDeadLetterPolicy">pulsar.ConsumerDeadLetterPolicy</a> | <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code> = <a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a>, </span><span class="rst-sig-param">crypto_failure_action: <code><a href="_pulsar.ConsumerCryptoFailureAction.html" class="internal-link" title="_pulsar.ConsumerCryptoFailureAction">ConsumerCryptoFailureAction</a></code> = ConsumerCryptoFailureAction.FAIL, </span><span class="rst-sig-param">is_pattern_topic: <code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code> = <a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a></span>) -&gt; <code><a href="pulsar.asyncio.Consumer.html" class="internal-link" title="pulsar.asyncio.Consumer">Consumer</a></code>:</span>
<a class="sourceLink" href="https://github.com/apache/pulsar-client-python/tree/v3.9.0/pulsar/asyncio.py#L548">
(source)
</a>
<a class="headerLink" href="#subscribe" title="pulsar.asyncio.Client.subscribe">
ΒΆ
</a>
</div>
<div class="docstring functionBody">
<div><p>Subscribe to the given topic and subscription combination.</p><table class="fieldTable"><tr class="fieldStart"><td class="fieldName" colspan="2">Parameters</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">topic:</span><code><a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a></code>, <code><a href="https://docs.python.org/3/library/typing.html#typing.List" class="intersphinx-link">List</a>[<a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a>]</code>, or regex pattern</td><td class="fieldArgDesc">The name of the topic, list of topics or regex pattern.
When <code>is_pattern_topic</code> is True, <code>topic</code> is treated as a regex.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">subscription<wbr></wbr>_name:</span><code><a href="https://docs.python.org/3/library/stdtypes.html#str" class="intersphinx-link">str</a></code></td><td class="fieldArgDesc">The name of the subscription.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">consumer<wbr></wbr>_type:</span><code><a href="_pulsar.ConsumerType.html" class="internal-link" title="_pulsar.ConsumerType">pulsar.ConsumerType</a></code>, <em>default</em> <code>pulsar.ConsumerType.Exclusive</code></td><td class="fieldArgDesc">Select the subscription type to be used when subscribing to the topic.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">schema:</span>pulsar.schema.Schema|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Define the schema of the data that will be received by this consumer.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">receiver<wbr></wbr>_queue<wbr></wbr>_size:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> <span class="literal">1000</span></td><td class="fieldArgDesc">Sets the size of the consumer receive queue.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">max<wbr></wbr>_total<wbr></wbr>_receiver<wbr></wbr>_queue<wbr></wbr>_size<wbr></wbr>_across<wbr></wbr>_partitions:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> <span class="literal">50000</span></td><td class="fieldArgDesc">Set the max total receiver queue size across partitions.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">consumer<wbr></wbr>_name:</span>str|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Sets the consumer name.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">unacked<wbr></wbr>_messages<wbr></wbr>_timeout<wbr></wbr>_ms:</span>int|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Sets the timeout in milliseconds for unacknowledged messages.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">broker<wbr></wbr>_consumer<wbr></wbr>_stats<wbr></wbr>_cache<wbr></wbr>_time<wbr></wbr>_ms:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> <span class="literal">30000</span></td><td class="fieldArgDesc">Sets the time duration for which the broker-side consumer stats
will be cached in the client.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">negative<wbr></wbr>_ack<wbr></wbr>_redelivery<wbr></wbr>_delay<wbr></wbr>_ms:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> <span class="literal">60000</span></td><td class="fieldArgDesc">The delay after which to redeliver the messages that failed to be
processed.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">is<wbr></wbr>_read<wbr></wbr>_compacted:</span><code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code>, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a></code></td><td class="fieldArgDesc">Selects whether to read the compacted version of the topic.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">properties:</span>dict|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Sets the properties for the consumer.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">initial<wbr></wbr>_position:</span><code><a href="_pulsar.InitialPosition.html" class="internal-link" title="_pulsar.InitialPosition">InitialPosition</a></code>, <em>default</em> <code>InitialPosition.Latest</code></td><td class="fieldArgDesc">Set the initial position of a consumer when subscribing to the topic.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">crypto<wbr></wbr>_key<wbr></wbr>_reader:</span>pulsar.CryptoKeyReader|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Symmetric encryption class implementation.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">replicate<wbr></wbr>_subscription<wbr></wbr>_state<wbr></wbr>_enabled:</span><code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code>, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a></code></td><td class="fieldArgDesc">Set whether the subscription status should be replicated.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">max<wbr></wbr>_pending<wbr></wbr>_chunked<wbr></wbr>_message:</span><code><a href="https://docs.python.org/3/library/functions.html#int" class="intersphinx-link">int</a></code>, <em>default</em> <span class="literal">10</span></td><td class="fieldArgDesc">Consumer buffers chunk messages into memory until it receives all the chunks.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">auto<wbr></wbr>_ack<wbr></wbr>_oldest<wbr></wbr>_chunked<wbr></wbr>_message<wbr></wbr>_on<wbr></wbr>_queue<wbr></wbr>_full:</span><code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code>, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a></code></td><td class="fieldArgDesc">Automatically acknowledge oldest chunked messages on queue
full.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">start<wbr></wbr>_message<wbr></wbr>_id<wbr></wbr>_inclusive:</span><code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code>, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a></code></td><td class="fieldArgDesc">Set the consumer to include the given position of any reset
operation.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">batch<wbr></wbr>_receive<wbr></wbr>_policy:</span>pulsar.ConsumerBatchReceivePolicy|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Set the batch collection policy for batch receiving.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">key<wbr></wbr>_shared<wbr></wbr>_policy:</span>pulsar.ConsumerKeySharedPolicy|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Set the key shared policy for use when the ConsumerType is
KeyShared.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">batch<wbr></wbr>_index<wbr></wbr>_ack<wbr></wbr>_enabled:</span><code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code>, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a></code></td><td class="fieldArgDesc">Enable the batch index acknowledgement.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">regex<wbr></wbr>_subscription<wbr></wbr>_mode:</span><code><a href="_pulsar.RegexSubscriptionMode.html" class="internal-link" title="_pulsar.RegexSubscriptionMode">RegexSubscriptionMode</a></code>, </td><td class="fieldArgDesc">default=RegexSubscriptionMode.PersistentOnly
Set the regex subscription mode for use when the topic is a regex
pattern.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">dead<wbr></wbr>_letter<wbr></wbr>_policy:</span>pulsar.ConsumerDeadLetterPolicy|<blockquote>
None</blockquote>
, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#None" class="intersphinx-link">None</a></code></td><td class="fieldArgDesc">Set dead letter policy for consumer.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">crypto<wbr></wbr>_failure<wbr></wbr>_action:</span><code><a href="_pulsar.ConsumerCryptoFailureAction.html" class="internal-link" title="_pulsar.ConsumerCryptoFailureAction">ConsumerCryptoFailureAction</a></code>, </td><td class="fieldArgDesc">default=ConsumerCryptoFailureAction.FAIL
Set the behavior when the decryption fails.</td></tr><tr><td class="fieldArgContainer"><span class="fieldArg">is<wbr></wbr>_pattern<wbr></wbr>_topic:</span><code><a href="https://docs.python.org/3/library/functions.html#bool" class="intersphinx-link">bool</a></code>, <em>default</em> <code><a href="https://docs.python.org/3/library/constants.html#False" class="intersphinx-link">False</a></code></td><td class="fieldArgDesc">Whether <code>topic</code> is a regex pattern. If it's True when <code>topic</code> is a list, a ValueError
will be raised.</td></tr></table><table class="fieldTable"><tr class="fieldStart"><td class="fieldName" colspan="2">Returns</td></tr><tr><td class="fieldArgContainer"><code><a href="pulsar.asyncio.Consumer.html" class="internal-link" title="pulsar.asyncio.Consumer">Consumer</a></code></td><td class="fieldArgDesc">The consumer created</td></tr></table><table class="fieldTable"><tr class="fieldStart"><td class="fieldName" colspan="2">Raises</td></tr><tr><td class="fieldArgContainer"><code><a href="pulsar.asyncio.PulsarException.html" class="internal-link" title="pulsar.asyncio.PulsarException">PulsarException</a></code></td><td></td></tr></table></div>
</div>
</div><div class="baseinstancevariable private">
<a name="pulsar.asyncio.Client._client">
</a>
<a name="_client">
</a>
<div class="functionHeader">
<span class="py-defname">_client</span>: <code><a href="_pulsar.Client.html" class="internal-link" title="_pulsar.Client">_pulsar.Client</a></code> =
<a class="sourceLink" href="https://github.com/apache/pulsar-client-python/tree/v3.9.0/pulsar/asyncio.py#L404">
(source)
</a>
<a class="headerLink" href="#_client" title="pulsar.asyncio.Client._client">
ΒΆ
</a>
</div>
<div class="functionBody">
<div><p class="undocumented">Undocumented</p></div>
</div>
</div>
</div>
</div>
</div>
<footer class="navbar navbar-default">
<hr />
<div class="container-fluid">
<a href="index.html">API Documentation</a> for pulsar/_pulsar,
generated by <a href="https://github.com/twisted/pydoctor/">pydoctor</a>
25.4.0 at 2025-12-30 16:47:40.
</div>
<script src="ajax.js" type="text/javascript"></script>
<script src="searchlib.js" type="text/javascript"></script>
<script src="search.js" type="text/javascript"></script>
</footer>
<script src="pydoctor.js" type="text/javascript"></script>
</body>
</html>