
Технология блокчейн и Apache Kafka имеют общие характеристики, которые предполагают естественную близость. Например, оба разделяют концепцию «неизменяемого журнала только с добавлением». В случае раздела Kafka:
Каждый раздел представляет собой упорядоченную неизменяемую последовательность записей, которая постоянно добавляется к структурированному журналу фиксации. Каждой записи в разделах присваивается последовательный идентификационный номер, называемый смещением, который однозначно идентифицирует каждую запись в разделе [Apache Kafka]
В то время как блокчейн можно описать как:
постоянно растущий список записей, называемых блоками, которые связаны и защищены с помощью криптографии. Каждый блок обычно содержит хеш-указатель в качестве ссылки на предыдущий блок, отметку времени и данные транзакции [Википедия]
Ясно, что эти технологии разделяют параллельные концепции неизменяемой последовательной структуры, причем Kafka особенно оптимизирован для высокой пропускной способности и горизонтальной масштабируемости, а блокчейн превосходит гарантию порядка и структуры последовательности.
Интегрируя эти технологии, мы можем создать платформу для экспериментов с концепциями блокчейн.
Kafka предоставляет удобную структуру для распределенного однорангового взаимодействия с некоторыми характеристиками, особенно подходящими для приложений блокчейна. Хотя этот подход может оказаться нежизнеспособным в ненадежной общедоступной среде, он может найти практическое применение в частной сети или сети консорциума. См. Масштабирование блокчейнов с помощью Apache Kafka для получения дополнительных идей о том, как это можно реализовать.
Кроме того, поэкспериментировав, мы сможем использовать концепции, уже реализованные в Kafka (например, сегментирование по разделам), чтобы изучить решения проблем блокчейна в общедоступных сетях (например, проблемы масштабируемости).
Поэтому цель этого эксперимента - взять простую реализацию блокчейна и перенести ее на платформу Kafka; мы возьмем концепцию последовательного журнала Кафки и гарантируем неизменность, связывая записи вместе с хешами. Тема blockchain на Kafka станет нашей распределенной бухгалтерской книгой. Графически это будет выглядеть так:

Введение в Кафку

Kafka - это потоковая платформа, предназначенная для высокопроизводительного обмена сообщениями в реальном времени, то есть она обеспечивает публикацию и подписку на потоки записей. В этом отношении он похож на очередь сообщений или традиционную корпоративную систему обмена сообщениями. Некоторые из характеристик:
- Высокая пропускная способность: брокеры Kafka могут поглощать гигабайты данных в секунду, переводя их в миллионы сообщений в секунду. Вы можете узнать больше о характеристиках масштабируемости в Бенчмаркинге Apache Kafka: 2 миллиона операций записи в секунду.
- Конкурирующие потребители: одновременная доставка сообщений нескольким потребителям, обычно дорогостоящая в традиционных системах обмена сообщениями, не сложнее, чем для одного потребителя. Это означает, что мы можем проектировать для конкурирующих потребителей, гарантируя, что каждый потребитель получит только одно из сообщений, и достигая высокой степени горизонтальной масштабируемости.
- Отказоустойчивость: репликация данных на нескольких узлах кластера сводит к минимуму влияние отказов отдельных узлов.
- Сохранение и воспроизведение сообщений: брокеры Kafka ведут учет зачетов потребителей - позиции потребителя в потоке сообщений. Используя это, потребители могут вернуться к предыдущей позиции в потоке, даже если сообщения уже были доставлены, что позволяет им воссоздать статус системы в определенный момент времени. Брокеры могут быть настроены на неограниченное хранение сообщений, что необходимо для приложений блокчейна.
В Kafka каждая тема разбита на разделы, где каждый раздел представляет собой последовательность записей, к которым постоянно добавляются. Это похоже на текстовый файл журнала, где в конец добавляются новые строки. Каждой записи в разделе присваивается последовательный идентификатор, называемый смещением, который однозначно идентифицирует запись.

Брокер Kafka может быть запрошен по смещению, то есть потребитель может сбросить свое смещение до некоторой произвольной точки в журнале, чтобы получить записи с этой точки вперед.
Руководство
Полный исходный код доступен здесь.
Предварительные требования
- Некоторое понимание концепций блокчейна: Учебное пособие, приведенное ниже, основано на реализациях от Daniel van Flymen и Gerald Nash, которые являются отличными практическими введениями. Следующее руководство в значительной степени основано на этих концепциях при использовании Kafka в качестве транспорта сообщений. Фактически, мы перенесем блокчейн Python в Kafka, сохранив большую часть текущей реализации.
- Базовые знания Python: код написан для Python 3.6.
- Docker: docker-compose используется для запуска брокера Kafka.
- Kafkacat: это полезный инструмент для взаимодействия с Kafka (например, публикация сообщений в темах).
При запуске наш потребитель Kafka попытается сделать три вещи: инициализировать новую цепочку блоков, если она еще не создана; построить внутреннее представление текущего состояния темы блокчейна; затем начните читать транзакции в цикле:
Шаг инициализации выглядит так:
Во-первых, мы находим максимально возможное смещение по теме блокчейна. Если в этой теме еще ничего не публиковалось, значит, блокчейн новый, поэтому мы начнем с создания и публикации блока генезиса:
В read_and_validate_chain() мы сначала создадим потребителя для чтения из темы blockchain:
Некоторые примечания к параметрам, с которыми мы создаем этого потребителя:
- Установка группы потребителей в
blockchaingroup позволяет брокеру сохранять ссылку на смещение, которого достигли потребители, для данного раздела и темы. auto_offset_reset=OffsetType.EARLIESTозначает, что мы начнем загрузку сообщений с начала темы.auto_commit_enable=Trueпериодически уведомляет брокера об только что полученном смещении (в отличие от фиксации вручную)reset_offset_on_start=True- это переключатель, который активируетauto_offset_resetдля потребителяconsumer_timeout_ms=5000заставит потребителя вернуться из метода через пять секунд, если новые сообщения не читаются (мы достигли конца цепочки)
Затем начинаем читать блочные сообщения из темы blockchain:
На каждое сообщение мы получаем:
- Если это первый блок в цепочке, пропустите проверку и добавьте к нашей внутренней копии (это генезисный блок)
- В противном случае проверьте, что блок действителен по отношению к предыдущему блоку, и добавьте его в нашу копию.
- Запомните смещение только что израсходованного блока.
В конце этого процесса мы загрузим всю цепочку, отбросив все недопустимые блоки, и у нас будет ссылка на смещение последнего блока.
На этом этапе мы готовы создать потребителя по теме transactions:
В нашем примере тема была создана с двумя разделами, чтобы продемонстрировать, как разделение работает в Kafka. Разделы настраиваются в файле docker-compose.yml, в этой строке:
KAFKA_CREATE_TOPICS=transactions:2:1,blockchain:1:1
transactions:2:1 указывает количество разделов и коэффициент репликации (т. Е. Сколько брокеров будут поддерживать копию данных в этом разделе).
На этот раз наш потребитель будет начинать с OffsetType.LATEST, поэтому мы получаем транзакции, опубликованные только с текущего времени.
Прикрепив потребителя к определенному разделу темы transactions, мы можем увеличить общую пропускную способность всех потребителей этой темы. Брокер Kafka будет равномерно распределять входящие сообщения по двум разделам темы транзакций, если мы не укажем раздел при публикации в теме. Это означает, что каждый потребитель будет отвечать за обработку 50% сообщений, что удвоит потенциальную пропускную способность отдельного потребителя.
Теперь мы можем начать принимать транзакции:
По мере поступления транзакций мы будем добавлять их во внутренний список. Каждые три транзакции мы будем создавать новый блок и вызывать mine():
- Сначала мы проверим, является ли наш блокчейн самым длинным в сети; является ли наше сохраненное смещение последним или другие узлы уже опубликовали более поздние блоки в цепочке блоков? Это наш консенсусный шаг.
- Если новые блоки уже были добавлены, мы будем использовать предыдущий
read_and_validate_chain, на этот раз предоставив последнее известное смещение, чтобы получить только новые блоки. - На этом этапе мы можем попытаться вычислить доказательство работы, основываясь на доказательстве из последнего блока.
- Чтобы вознаградить себя за решение доказательства работы, мы можем вставить транзакцию в блок, выплачивая себе небольшое вознаграждение за блок.
- Наконец, мы опубликуем наш блок в теме блокчейна. Метод публикации выглядит так:
В действии
- Сначала запустите брокера:
docker-compose up -d
2. Запустите потребителя на разделе 0:
python kafka_blockchain.py 0
3. Опубликуйте 3 транзакции непосредственно в разделе 0:
4. Убедитесь, что транзакции добавлены в блок по теме blockchain:
kafkacat -C -b kafka:9092 -t blockchain
Вы должны увидеть такой результат:
Чтобы сбалансировать транзакции между двумя потребителями, запустите второго потребителя в разделе 1 и удалите -p 0 из сценария публикации выше.
Заключение
Kafka может обеспечить основу для простой платформы для экспериментов с блокчейном. Мы можем воспользоваться функциями, встроенными в платформу, и связанными с ними инструментами, такими как kafkacat, для экспериментов с распределенными одноранговыми транзакциями.
В то время как масштабирование транзакций в общедоступной среде представляет собой один набор проблем, в частной сети или консорциуме, где уже установлено реальное доверие, масштабирование транзакций может быть достигнуто с помощью реализации, которая использует преимущества концепции Kafka.