Skip to main content

async-kinesis-client

Python Kinesis Client library utilising asyncio

Based on Kinesis-Python project by Evan Borgstrom eborgstrom@nerdwallet.com https://github.com/NerdWalletOSS/kinesis-python but with asyncio magic

The problem with Kinesis-Python is that all the data ends up in a single thread and being checkpointed from there - so despite having many processes, the client is clogged by checkpointing. Besides, it checkpoints every single record and this is not configurable.

This client is based on aioboto3 library and uses Python 3.6+ async methods.

Usage:

import asyncio
from async_kinesis_client.kinesis_consumer import AsyncKinesisConsumer

async def read_stream():

    # This is a coroutine that reads all the records from a shard
    async def read_records(shard_reader):
        async for records in shard_reader.get_records():
            for r in records:
                print('Shard: {}; Record: {}'.format(shard_reader.shard_id, r))

    consumer = AsyncKinesisConsumer(
                stream_name='my-stream',
                checkpoint_table='my-checkpoint-table')

    # consumer will yield existing shards and will continue yielding
    # new shards if re-sharding happens             
    async for shard_reader in consumer.get_shard_readers():
        print('Got shard reader for shard id: {}'.format(shard_reader.shard_id))
        asyncio.ensure_future(read_records(shard_reader)) 

asyncio.get_event_loop().run_until_complete(read_stream())

AsyncShardReader and AsyncKinesisConsumer can be stopped from parallel coroutine by calling stop() method, consumer will stop all shard readers in that case. If you want to be notified of shard closing, catch ShardClosedException while reading records

AsyncShardReader exposes property millis_behind_latest which could be useful for determining application performance.

AsyncKinesisConsumer has following configuration methods:

set_checkpoint_interval(records) - how many records to skip before checkpointing

set_lock_duration(time) - how many seconds to hold the lock. Consumer would attempt to refresh the lock before that time

set_reader_sleep_time(time) - how long should shard reader wait (in seconds, fractions possible) if it did not receive any records from Kinesis stream

set_checkpoint_callback(coro) - set callback coroutine to be called before checkpointing next batch of records. Coroutine arguments: ShardId, SequenceNumber

Producer is rather trivial:

from async_kinesis_client.kinesis_producer import AsyncKinesisProducer

# ...

async def write_stream():
    producer = AsyncKinesisProducer(
        stream_name='my-stream',
        ordered=True
    )

    await producer.put_record(
        record=b'bytes',
        partition_key='string',     # optional, if none, default time-based key is used
        explicit_hash_key='string'  # optional
    )

Sending multiple records at once:

from async_kinesis_client.kinesis_producer import AsyncKinesisProducer

# ...

async def write_stream():
    producer = AsyncKinesisProducer(
        stream_name='my-stream',
        ordered=True
    )

    records = [
        {
            'Data': b'bytes',
            'PartitionKey': 'string',   # optional, if none, default time-based key is used
            'ExplicitHashKey': 'string' # optional
        },
        ...
    ]

    response = await producer.put_records(
        records=records
    )

    # See boto3 docs for response structure:
    # https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/kinesis.html#Kinesis.Client.put_records

AWS authentication. For testing outside AWS cloud, especially when Mutil-Factor Authentication is in use I find following snippet extremely useful:

import os
import aioboto3
from botocore import credentials
from aiobotocore import AioSession

    working_dir = os.path.join(os.path.expanduser('~'), '.aws/cli/cache')
    session = AioSession(profile=os.environ.get('AWS_PROFILE'))
    provider = session.get_component('credential_provider').get_provider('assume-role')
    provider.cache = credentials.JSONFileCache(working_dir)
    aioboto3.setup_default_session(botocore_session=session)

This allows re-using cached session token after completing any aws command under awsudo, all you need is to set AWS_PROFILE environment variable.

Currently library still not tested enough for different network events. Use it on your own risk, you've been warned.

Release files for async-kinesis-client 0.2.14

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for async-kinesis-client 0.2.14
File Size Uploaded
async-kinesis-client-0.2.14.tar.gz 14.0 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for async-kinesis-client 0.2.14
File Interpreter ABI Platform
async_kinesis_client-0.2.14-py3-none-any.whl Python 3 none any Details

Total release size: 29.4 kB

Release files / async-kinesis-client-0.2.14.tar.gz

Download URL async-kinesis-client-0.2.14.tar.gz
Size 14.0 kB
Tags Source
SHA-256 checksum
How to use checksums
1908928e432562c9c3e5e9c37224ff6f587345058732ae37beb18f6405bca562
BLAKE2b-256 checksum
How to use checksums
b05ee2d9b293c8cd76605beba210241b54badc8e33c1e46dd584a5681e0c0247
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/1.12.1 pkginfo/1.4.2 requests/2.20.1 setuptools/41.0.1 requests-toolbelt/0.8.0 tqdm/4.28.1 CPython/3.7.4

Release files / async_kinesis_client-0.2.14-py3-none-any.whl

Download URL async_kinesis_client-0.2.14-py3-none-any.whl
Size 15.4 kB
Tags Python 3
SHA-256 checksum
How to use checksums
6c7031226fcf531c7b0316dead11e3e0fb141cfcf477927e6d3bfb6ede21d0bc
BLAKE2b-256 checksum
How to use checksums
d0cee8653cbe0468faeb1997514e695e05896efa1bebc4cb15f9b9b4b24a71f3
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/1.12.1 pkginfo/1.4.2 requests/2.20.1 setuptools/41.0.1 requests-toolbelt/0.8.0 tqdm/4.28.1 CPython/3.7.4

Release history Release notifications | RSS feed

This release

0.2.14 This release

2 release files

0.2.12

2 release files

0.2.11

2 release files

0.2.10

2 release files

0.2.9

2 release files

0.2.8

2 release files

0.2.7

2 release files

0.2.6

2 release files

0.2.5

2 release files

0.2.4

2 release files

0.2.3

2 release files

0.2.2

2 release files

0.2.1

2 release files

0.2.0

2 release files

0.1.7

2 release files

0.1.6

2 release files

0.1.5

2 release files

0.1.4

2 release files

0.1.3

2 release files

0.1.2

2 release files

0.1.1

2 release files

0.1.0

2 release files

0.0.4

2 release files

0.0.3

2 release files

0.0.2

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page