Yandex Cloud
Поиск
Связаться с экспертомПопробовать бесплатно
  • Кейсы
  • Документация
  • Блог
  • Все сервисы
    • Cloud Interconnect
    • Cloud Backup
    • Cloud Registry
    • Yandex AI Studio
    • Compute Cloud
    • Object Storage
    • Managed Service for Kubernetes®
    • Yandex BareMetal
    • Smart Web Security
    • Security Deck
    • Managed Service for PostgreSQL
    • Managed Service for ClickHouse®
    • Monium
    • Cloud CDN
    • Network Load Balancer
    • Virtual Private Cloud
    • Cloud DNS
    • Application Load Balancer
    • Yandex Cloud Video
    • Stackland
    • Yandex Cloud Router
    • Yandex Managed Service for Trino
    • Managed Service for MySQL®
    • Managed Service for Valkey™
    • Managed Service for Apache Spark™
    • Yandex StoreDoc
    • Managed Service for OpenSearch
    • Managed Service for Apache Kafka®
    • Data Transfer
    • Yandex MPP Analytics Engine for PostgreSQL
    • Yandex Managed Service for Apache Airflow®
    • Data Processing
    • Yandex MetaData Hub
    • Managed Service for YDB
    • Managed Service for Sharded PostgreSQL
    • Managed Service for YTsaurus
    • Yandex WebSQL
    • DataLens
    • Yandex Search API
    • SpeechSense
    • SpeechKit
    • DataSphere
    • Vision OCR
    • Translate
    • Yandex Neurosupport
    • Yandex Cloud Detection and Response
    • Yandex Identity Hub
    • Key Management Service
    • Certificate Manager
    • Yandex Lockbox
    • Audit Trails
    • SmartCaptcha
    • Cloud Desktop
    • GOST Gateway
    • Yandex SIEM
    • SourceCraft Code Assistant
    • Container Registry
    • Managed Service for GitLab
    • SourceCraft
    • Managed Service for Prometheus®
    • Cloud Functions
    • API Gateway
    • Yandex Cloud Postbox
    • Message Queue
    • Serverless Integrations
    • IoT Core
    • Data Streams
    • Serverless Containers
    • Cloud Notification Service
    • Yandex Query
    • Identity and Access Management
    • Yandex Cloud Console
    • Resource Manager
    • Yandex Cloud Billing
    • Yandex Cloud Quota Manager
    • Cloud Apps
  • Статус работы сервисов
  • Marketplace
    • Популярные
    • Инфраструктура и сеть
    • Платформа данных
    • Искусственный интеллект
    • Безопасность
    • Инструменты DevOps
    • Бессерверные вычисления
    • Управление ресурсами
  • Все решения
    • По отраслям
    • По типу задач
    • Экономика платформы
    • Yandex Cloud Trust
    • Техническая поддержка
    • Каталог партнёров
    • Обучение и сертификация
    • Облако для стартапов
    • Облако для крупного бизнеса
    • Центр технологий для общества
    • Облако для интеграторов
    • Поддержка IT-бизнеса
    • Облако для фрилансеров
    • Обучение и сертификация
    • Блог
    • Документация
    • Контент-программа
    • Мероприятия и вебинары
    • Контакты, чаты и сообщества
    • Идеи
    • Калькулятор цен
    • Тарифы
    • Акции и free tier
  • Кейсы
  • Документация
  • Блог
Создавайте контент и получайте гранты!Готовы написать своё руководство? Участвуйте в контент-программе и получайте гранты на работу с облачными сервисами!
Подробнее о программе
Проект Яндекса
© 2026 ООО «Яндекс.Облако»
Yandex Data Streams
    • Все инструкции
    • Управление потоками данных
      • Подготовка окружения
      • Создание потока данных
      • Отправка данных в поток
      • Чтение данных из потока
      • Удаление потока данных
  • Управление доступом
  • Правила тарификации
  • Вопросы и ответы
  1. Пошаговые инструкции
  2. Работа с AWS SDK
  3. Чтение данных из потока

Чтение данных из потока в AWS SDK

Статья создана
Yandex Cloud
Обновлена 28 декабря 2023 г.
Открыть в Markdown
Python

Для чтения записей из потока данных используется пара методов: get_shard_iterator и get_record/get_records. При вызове этого метода необходимо указать следующие параметры:

  • Имя потока данных, например example-stream.
  • Идентификатор облака, в котором находится поток, например b1gi1kuj2dht********.
  • Идентификатор базы данных YDB с потоком, например cc8028jgtuab********.

Вам также потребуется настроить AWS SDK и назначить сервисному аккаунту роль yds.viewer.

Для чтения записей из потока с параметрами, указанными выше:

  1. Создайте файл stream_get_records.py и скопируйте в него следующий код:

    import boto3
    from pprint import pprint
    import itertools
    
    def get_records(cloud, database, stream_name):
        client = boto3.client('kinesis', endpoint_url="https://yds.serverless.yandexcloud.net")
    
        StreamName = "/ru-central1/{cloud}/{database}/{stream}".format(cloud=cloud,
                                                                     database=database,
                                                                     stream=stream_name)
    
    
        describe_stream_result = client.describe_stream(StreamName=StreamName)
        shard_iterators = {}
    
        shards = [shard["ShardId"] for shard in describe_stream_result['StreamDescription']['Shards']]
    
        for shard_id in itertools.cycle(shards):
            if shard_id not in shard_iterators:
                shard_iterators[shard_id] = client.get_shard_iterator(StreamName=StreamName,
                                                                     ShardId=shard_id,
                                                                     ShardIteratorType='LATEST')['ShardIterator']
               
            record_response = client.get_records(ShardIterator=shard_iterators[shard_id])
            if "Records" in record_response:
                for record in [record for record in record_response["Records"]]:
                    yield record["Data"]
    
            if "NextShardIterator" in record_response:
                shard_iterators[shard_id] = record_response["NextShardIterator"]
    
    
    if __name__ == '__main__':
        for record in get_records(cloud="b1gi1kuj2dht********",
                                  database="cc8028jgtuab********",
                                  stream_name="example-stream"):
            pprint(record)    
            print("The record has been read successfully")
    
  2. Запустите программу:

    python3 stream_get_records.py
    

    Результат:

    The record has been read successfully
    b'{"user_id":"user1","score":100}'
    The record has been read successfully
    b'{"user_id":"user1","score":100}'
    ...
    

Была ли статья полезна?

Предыдущая
Отправка данных в поток
Следующая
Удаление потока данных
Создавайте контент и получайте гранты!Готовы написать своё руководство? Участвуйте в контент-программе и получайте гранты на работу с облачными сервисами!
Подробнее о программе
Проект Яндекса
© 2026 ООО «Яндекс.Облако»