Быстрый старт: Сбор данных Event Hubs в службе хранилища Azure и чтение с его помощью с помощью Python (azure-eventhub)

Вы можете настроить концентратор событий, чтобы данные, отправленные в концентратор событий, были записаны в учетной записи хранения Azure или Azure Data Lake Storage 1-го поколения или 2-го поколения. В этой статье показано, как написать код Python для отправки событий в концентратор событий и считывать собранные данные из хранилища BLOB-объектов Azure. Дополнительные сведения об этой функции см. в обзоре функции "Сбор центров событий".

В этом кратком руководстве используется пакет SDK для Python Azure для демонстрации функции отслеживания. Приложение sender.py отправляет имитацию телеметрии окружающей среды в центры событий в формате JSON. Концентратор событий настроен для использования функции Capture для записи этих данных в хранилище BLOB-объектов в пакетах. Приложение capturereader.py считывает эти двоичные объекты и создает файл дополнения для каждого устройства. Затем приложение записывает данные в CSV-файлы.

В этом кратком руководстве вы сможете:

  • Создайте учетную запись хранения BLOB-объектов Azure и контейнер на портале Azure.
  • Создайте пространство имен Центров событий с помощью портала Azure.
  • Создайте концентратор событий с включенной функцией записи и подключите ее к учетной записи хранения.
  • Отправка данных в концентратор событий с помощью скрипта Python.
  • Чтение и обработка файлов из Event Hubs Capture с использованием другого скрипта на Python.

Предпосылки

Включение функции записи для концентратора событий

Включите функцию записи для концентратора событий. Для этого следуйте инструкциям в разделе Включение функции Capture в Центрах событий через портал Azure. Выберите учетную запись службы хранения и контейнер для BLOB-объектов, созданный на предыдущем шаге. Выберите avro для формата сериализации событий вывода.

Создание скрипта Python для отправки событий в концентратор событий

В этом разделе описано, как создать скрипт Python, который отправляет 200 событий (10 устройств * 20 событий) в концентратор событий. Эти события являются примером чтения среды, которое отправляется в формате JSON.

  1. Откройте любимый редактор Python, например Visual Studio Code.

  2. Создайте скрипт с именем sender.py.

  3. Вставьте следующий код в sender.py.

    import time
    import os
    import uuid
    import datetime
    import random
    import json
    
    from azure.eventhub import EventHubProducerClient, EventData
    
    # This script simulates the production of events for 10 devices.
    devices = []
    for x in range(0, 10):
        devices.append(str(uuid.uuid4()))
    
    # Create a producer client to produce and publish events to the event hub.
    producer = EventHubProducerClient.from_connection_string(conn_str="EVENT HUBS NAMESPACE CONNECTION STRING", eventhub_name="EVENT HUB NAME")
    
    for y in range(0,20):    # For each device, produce 20 events. 
        event_data_batch = producer.create_batch() # Create a batch. You will add events to the batch later. 
        for dev in devices:
            # Create a dummy reading.
        reading = {
                'id': dev, 
                'timestamp': str(datetime.datetime.utcnow()), 
                'uv': random.random(), 
                'temperature': random.randint(70, 100), 
                'humidity': random.randint(70, 100)
            }
            s = json.dumps(reading) # Convert the reading into a JSON string.
            event_data_batch.add(EventData(s)) # Add event data to the batch.
        producer.send_batch(event_data_batch) # Send the batch of events to the event hub.
    
    # Close the producer.    
    producer.close()
    
  4. Замените следующие значения в скриптах:

    • Замените EVENT HUBS NAMESPACE CONNECTION STRING на строку подключения для пространства имен центров событий.
    • Замените EVENT HUB NAME именем концентратора событий.
  5. Запустите скрипт, чтобы отправить события в концентратор событий.

  6. На портале Azure можно убедиться, что концентратор событий получил сообщения. Перейдите в представление "Сообщения" в разделе "Метрики ". Обновите страницу, чтобы обновить диаграмму. Для отображения полученных сообщений может потребоваться несколько секунд.

    Убедитесь, что концентратор событий получил сообщения

Создайте скрипт на Python для чтения ваших Capture файлов

В этом примере собранные данные хранятся в хранилище BLOB-объектов Azure. Скрипт в этом разделе считывает захваченные файлы данных из учетной записи хранения Azure и создает CSV-файлы, которые можно легко открыть и просмотреть. Вы увидите 10 файлов в текущем рабочем каталоге приложения. Эти файлы содержат показания среды для 10 устройств.

  1. В редакторе Python создайте скрипт с именем capturereader.py. Этот скрипт считывает захваченные файлы и создает файл для каждого устройства для записи данных только для этого устройства.

  2. Вставьте следующий код в capturereader.py.

    import os
    import string
    import json
    import uuid
    import avro.schema
    
    from azure.storage.blob import ContainerClient, BlobClient
    from avro.datafile import DataFileReader, DataFileWriter
    from avro.io import DatumReader, DatumWriter
    
    
    def processBlob2(filename):
        reader = DataFileReader(open(filename, 'rb'), DatumReader())
        dict = {}
        for reading in reader:
            parsed_json = json.loads(reading["Body"])
            if not 'id' in parsed_json:
                return
            if not parsed_json['id'] in dict:
                list = []
                dict[parsed_json['id']] = list
            else:
                list = dict[parsed_json['id']]
                list.append(parsed_json)
        reader.close()
        for device in dict.keys():
            filename = os.getcwd() + '\\' + str(device) + '.csv'
            deviceFile = open(filename, "a")
            for r in dict[device]:
                deviceFile.write(", ".join([str(r[x]) for x in r.keys()])+'\n')
    
    def startProcessing():
        print('Processor started using path: ' + os.getcwd())
        # Create a blob container client.
        container = ContainerClient.from_connection_string("AZURE STORAGE CONNECTION STRING", container_name="BLOB CONTAINER NAME")
        blob_list = container.list_blobs() # List all the blobs in the container.
        for blob in blob_list:
            # Content_length == 508 is an empty file, so process only content_length > 508 (skip empty files).        
            if blob.size > 508:
                print('Downloaded a non empty blob: ' + blob.name)
                # Create a blob client for the blob.
                blob_client = ContainerClient.get_blob_client(container, blob=blob.name)
                # Construct a file name based on the blob name.
                cleanName = str.replace(blob.name, '/', '_')
                cleanName = os.getcwd() + '\\' + cleanName 
                with open(cleanName, "wb+") as my_file: # Open the file to write. Create it if it doesn't exist. 
                    my_file.write(blob_client.download_blob().readall()) # Write blob contents into the file.
                processBlob2(cleanName) # Convert the file into a CSV file.
                os.remove(cleanName) # Remove the original downloaded file.
                # Delete the blob from the container after it's read.
                container.delete_blob(blob.name)
    
    startProcessing()    
    
  3. Замените строку подключения AZURE STORAGE CONNECTION STRING вашей учетной записи Azure Storage. Имя контейнера, созданного в этом кратком руководстве, — capture. Если для контейнера использовалось другое имя, замените запись именем контейнера в учетной записи хранения.

Запуск скриптов

  1. Откройте командную строку с Python в пути, а затем выполните следующие команды, чтобы установить пакеты необходимых компонентов Python:

    pip install azure-storage-blob
    pip install azure-eventhub
    pip install avro-python3
    
  2. Измените каталог на каталог, в котором вы сохранили sender.py и capturereader.py, и выполните следующую команду:

    python sender.py
    

    Эта команда запускает новый процесс Python для запуска отправителя.

  3. Подождите несколько минут, пока запись будет запущена, а затем введите следующую команду в исходном окне команды:

    python capturereader.py
    

    Этот обработчик захвата использует локальный каталог для скачивания всех больших двоичных объектов из контейнера и учетной записи хранения. Он обрабатывает файлы, которые не пусты, и записывает результаты в виде CSV-файлов в локальный каталог.

Дальнейшие шаги

Ознакомьтесь с примерами Python на сайте GitHub.