Встреча создана!
Бот @JFB_admin_bot выдаст ссылку. Для корректной работы лучше его перезапустить)
Залетаем✈️
Бот @JFB_admin_bot выдаст ссылку. Для корректной работы лучше его перезапустить)
Залетаем
Please open Telegram to view this post
VIEW IN TELEGRAM
🔥3
Tree в Java.
Самая сложная коллекция.
В представленном видео, мы подробно разобрали, что такое Tree и как они представлены в Java.
Огромное спасибо @ElizaFanat за рассказ и демонстрации.🙂
Ссылка на Youtube
Ссылка на Рутьюб
Смотрите, ставьте лайки, подписывайтесь на каналы!✌️
Самая сложная коллекция.
В представленном видео, мы подробно разобрали, что такое Tree и как они представлены в Java.
Огромное спасибо @ElizaFanat за рассказ и демонстрации.
Ссылка на Youtube
Ссылка на Рутьюб
Смотрите, ставьте лайки, подписывайтесь на каналы!
Please open Telegram to view this post
VIEW IN TELEGRAM
👍3🔥2
История IT-технологий сегодня — 29 декабря
ℹ️ Кто родился в этот день
Ян Ливингстон (англ. Ian Livingstone; 29 декабря 1949, Престбери, Чешир, Англия) — английский автор-фантаст и антрепренёр. Соавтор первой книги-игры The Warlock of Firetop Mountain из серии Fighting Fantasy и сооснователь Games Workshop.
🌐 Знаковые события
Не нашел(
#Biography #Birth_Date #Events #29Декабря
Ян Ливингстон (англ. Ian Livingstone; 29 декабря 1949, Престбери, Чешир, Англия) — английский автор-фантаст и антрепренёр. Соавтор первой книги-игры The Warlock of Firetop Mountain из серии Fighting Fantasy и сооснователь Games Workshop.
Не нашел(
#Biography #Birth_Date #Events #29Декабря
Please open Telegram to view this post
VIEW IN TELEGRAM
👍2
Глубокая архитектура и внутреннее устройство RabbitMQ
AMQP 1.0 vs AMQP 0-9-1: эволюция протокола
AMQP 0-9-1: классическая модель RabbitMQ
AMQP 0-9-1 (Advanced Message Queuing Protocol версии 0-9-1) — это протокол, вокруг которого строился RabbitMQ с момента его создания.
Его ключевые характеристики:
Строгая топологическая модель с явным объявлением exchanges, queues и bindings
Frame-based протокол с четкой структурой кадров (frames)
Каналы (Channels) — виртуальные соединения внутри одного TCP-соединения для уменьшения накладных расходов
Подтверждения (Acknowledgements) на уровне потребителя и издателя
Транзакции через механизм tx.commit/tx.rollback
В 2025 году AMQP 0-9-1 остается основным протоколом для RabbitMQ, особенно для сценариев, требующих сложной маршрутизации и гарантий доставки.
AMQP 1.0: стандартизация и упрощение
AMQP 1.0 — это стандартизированная версия протокола, разработанная OASIS, с существенными изменениями:
Более абстрактная модель без явных понятий exchanges и bindings
Сообщения как первоклассные сущности с расширенными заголовками и свойствами
Связи (Links) вместо каналов, с разделением на sender и receiver links
Улучшенная обработка ошибок со стандартизированными кодами
Поддержка транзакций через распределенные транзакции (не полностью реализовано в RabbitMQ)
Сравнительный анализ
Когда использовать AMQP 0-9-1:
При работе со сложными маршрутизационными сценариями (headers exchange, topic exchange)
Когда необходима максимальная совместимость с существующими клиентскими библиотеками
Для использования расширенных функций RabbitMQ (политики, shovel, federation)
При миграции legacy систем без переписывания логики маршрутизации
Когда использовать AMQP 1.0:
При интеграции с системами, поддерживающими только AMQP 1.0 (Azure Service Bus, Apache Qpid)
Для межплатформенной совместимости в гетерогенных средах
Когда требуется стандартизированная обработка сообщений без vendor lock-in
В сценариях, где важнее семантика сообщений, а не топология маршрутизации
Поддержка в RabbitMQ 3.13+
RabbitMQ поддерживает оба протокола одновременно через разные порты:
AMQP 0-9-1: порт 5672 по умолчанию
AMQP 1.0: порт 5671 (или отдельно настроенный)
Важно понимать, что это два различных протокола с разными клиентскими библиотеками и семантикой. RabbitMQ выступает в роли моста между ними, но с ограничениями в преобразовании расширенных функций.
Низкоуровневая архитектура: как работает RabbitMQ внутри
Основа на Erlang/OTP: философия отказоустойчивости
Erlang/OTP — это платформа для построения распределенных, отказоустойчивых систем с soft real-time характеристиками.
Выбор Erlang для RabbitMQ был стратегическим решением:
Actor model — каждый процесс Erlang (не системный процесс) изолирован и обрабатывает сообщения асинхронно
"Let it crash" философия — процессы проектируются с ожиданием сбоев, которые обрабатываются супервизорами
Hot code reloading — возможность обновления кода без остановки системы
Распределенная природа — встроенная поддержка кластеризации и межпроцессного взаимодействия
В RabbitMQ различные компоненты реализованы как процессы Erlang:
Одна очередь = один или несколько процессов Erlang
Каждое соединение клиента = процесс Erlang
Каждый канал внутри соединения = отдельный процесс
#Java #middle #RabbitMQ
AMQP 1.0 vs AMQP 0-9-1: эволюция протокола
AMQP 0-9-1: классическая модель RabbitMQ
AMQP 0-9-1 (Advanced Message Queuing Protocol версии 0-9-1) — это протокол, вокруг которого строился RabbitMQ с момента его создания.
Его ключевые характеристики:
Строгая топологическая модель с явным объявлением exchanges, queues и bindings
Frame-based протокол с четкой структурой кадров (frames)
Каналы (Channels) — виртуальные соединения внутри одного TCP-соединения для уменьшения накладных расходов
Подтверждения (Acknowledgements) на уровне потребителя и издателя
Транзакции через механизм tx.commit/tx.rollback
В 2025 году AMQP 0-9-1 остается основным протоколом для RabbitMQ, особенно для сценариев, требующих сложной маршрутизации и гарантий доставки.
AMQP 1.0: стандартизация и упрощение
AMQP 1.0 — это стандартизированная версия протокола, разработанная OASIS, с существенными изменениями:
Более абстрактная модель без явных понятий exchanges и bindings
Сообщения как первоклассные сущности с расширенными заголовками и свойствами
Связи (Links) вместо каналов, с разделением на sender и receiver links
Улучшенная обработка ошибок со стандартизированными кодами
Поддержка транзакций через распределенные транзакции (не полностью реализовано в RabbitMQ)
Сравнительный анализ
Когда использовать AMQP 0-9-1:
При работе со сложными маршрутизационными сценариями (headers exchange, topic exchange)
Когда необходима максимальная совместимость с существующими клиентскими библиотеками
Для использования расширенных функций RabbitMQ (политики, shovel, federation)
При миграции legacy систем без переписывания логики маршрутизации
Когда использовать AMQP 1.0:
При интеграции с системами, поддерживающими только AMQP 1.0 (Azure Service Bus, Apache Qpid)
Для межплатформенной совместимости в гетерогенных средах
Когда требуется стандартизированная обработка сообщений без vendor lock-in
В сценариях, где важнее семантика сообщений, а не топология маршрутизации
Поддержка в RabbitMQ 3.13+
RabbitMQ поддерживает оба протокола одновременно через разные порты:
AMQP 0-9-1: порт 5672 по умолчанию
AMQP 1.0: порт 5671 (или отдельно настроенный)
Важно понимать, что это два различных протокола с разными клиентскими библиотеками и семантикой. RabbitMQ выступает в роли моста между ними, но с ограничениями в преобразовании расширенных функций.
Низкоуровневая архитектура: как работает RabbitMQ внутри
Основа на Erlang/OTP: философия отказоустойчивости
Erlang/OTP — это платформа для построения распределенных, отказоустойчивых систем с soft real-time характеристиками.
Выбор Erlang для RabbitMQ был стратегическим решением:
Actor model — каждый процесс Erlang (не системный процесс) изолирован и обрабатывает сообщения асинхронно
"Let it crash" философия — процессы проектируются с ожиданием сбоев, которые обрабатываются супервизорами
Hot code reloading — возможность обновления кода без остановки системы
Распределенная природа — встроенная поддержка кластеризации и межпроцессного взаимодействия
В RabbitMQ различные компоненты реализованы как процессы Erlang:
Одна очередь = один или несколько процессов Erlang
Каждое соединение клиента = процесс Erlang
Каждый канал внутри соединения = отдельный процесс
#Java #middle #RabbitMQ
👍2
Процессная модель и планировщики
Erlang VM использует вытесняющую многозадачность с планировщиками (schedulers).
В RabbitMQ 3.13+:
По умолчанию используется по одному планировщику на CPU core
Каждый планировщик имеет свою очередь исполняемых процессов
Reductions — единица измерения работы в Erlang, используемая для fair scheduling
Псевдокод работы планировщика Erlang:
Управление памятью и сборка мусора
Память в Erlang управляется через per-process heap и shared binary heap:
Process heap — небольшая частная куча для термов Erlang
Binary heap — общая куча для больших данных (тела сообщений в RabbitMQ)
Copying garbage collector для process heap, работающий при заполнении кучи
Reference counting для binary heap
Для сообщений размером более 64 байт (настраиваемый параметр) тело хранится в binary heap, а в очереди сохраняется только ссылка. Это позволяет эффективно обрабатывать большие сообщения с несколькими потребителями.
Protocol internals: от байтов к семантике
AMQP 0-9-1 Frame структура
Пример последовательности фреймов для публикации:
Процесс обработки входящего сообщения
#Java #middle #RabbitMQ
Erlang VM использует вытесняющую многозадачность с планировщиками (schedulers).
В RabbitMQ 3.13+:
По умолчанию используется по одному планировщику на CPU core
Каждый планировщик имеет свою очередь исполняемых процессов
Reductions — единица измерения работы в Erlang, используемая для fair scheduling
Псевдокод работы планировщика Erlang:
function scheduler_loop(queue, time_slice) {
while (true) {
process = queue.dequeue()
reductions_executed = 0
while (process.has_messages() && reductions_executed < time_slice) {
message = process.next_message()
result = process.execute(message)
reductions_executed += calculate_reductions(result)
if (process.crashed()) {
notify_supervisor(process)
break
}
}
if (process.has_messages()) {
queue.enqueue(process) // Вернуть в конец очереди
}
}
}Управление памятью и сборка мусора
Память в Erlang управляется через per-process heap и shared binary heap:
Process heap — небольшая частная куча для термов Erlang
Binary heap — общая куча для больших данных (тела сообщений в RabbitMQ)
Copying garbage collector для process heap, работающий при заполнении кучи
Reference counting для binary heap
Для сообщений размером более 64 байт (настраиваемый параметр) тело хранится в binary heap, а в очереди сохраняется только ссылка. Это позволяет эффективно обрабатывать большие сообщения с несколькими потребителями.
Protocol internals: от байтов к семантике
AMQP 0-9-1 Frame структура
Frame Structure:
+----------+----------+----------+----------+----------+----------+
| Type | Channel | Size | Payload | Frame End |
| (1 byte) | (2 bytes)| (4 bytes)| (size bytes) | (1 byte) |
+----------+----------+----------+----------+----------+----------+
Frame Types:
- METHOD (1) : Вызов метода AMQP (declare, publish, consume)
- HEADER (2) : Заголовки сообщения (properties, headers)
- BODY (3) : Часть тела сообщения (может быть несколько фреймов)
- HEARTBEAT (8) : Keep-alive фрейм
Пример последовательности фреймов для публикации:
[METHOD] Basic.Publish(exchange="amq.direct", routing_key="queue1")
[HEADER] properties={content_type: "text/plain"}, body_size=1024
[BODY] chunk 1 of 1024 bytes
[BODY] chunk 2 of 1024 bytes (если сообщение больше frame_max)
Процесс обработки входящего сообщения
// Псевдокод обработки publish на стороне брокера
handle_publish(frame) {
// 1. Парсинг и валидация фреймов
method_frame = parse_method_frame(frame)
header_frame = read_next_frame() // Ожидаем HEADER фрейм
// 2. Поиск exchange по имени
exchange = lookup_exchange(method_frame.exchange)
if (!exchange) {
if (method_frame.mandatory) {
send_basic_return() // Сообщение возвращается отправителю
}
return
}
// 3. Маршрутизация через exchange
routes = exchange.route(method_frame.routing_key, header_frame.properties)
// 4. Для каждого получателя (очереди)
for (queue in routes.queues) {
// 5. Проверка TTL сообщения
if (header_frame.properties.expiration && is_expired(header_frame)) {
continue
}
// 6. Сохранение в очередь
message = {
id: generate_message_id(),
properties: header_frame.properties,
body: read_body_frames() // Чтение всех BODY фреймов
}
// 7. В зависимости от типа очереди
if (queue.type == "quorum") {
quorum_queue_append(queue, message)
} else if (queue.type == "stream") {
stream_append(queue, message)
} else {
classic_queue_append(queue, message)
}
}
// 8. Подтверждение издателю (если включены publisher confirms)
if (connection.publisher_confirms) {
send_basic_ack(delivery_tag)
}
}
#Java #middle #RabbitMQ
👍2
Новые типы очередей в RabbitMQ 3.13+
Quorum Queues: консенсус как основа надежности
Quorum Queues используют алгоритм Raft для репликации данных между узлами кластера. Raft — это алгоритм консенсуса, обеспечивающий согласованное состояние распределенной системы при наличии отказов.
Архитектура Quorum Queue
Жизненный цикл сообщения в Quorum Queue
Преимущества Quorum Queues:
Автоматическое восстановление после потери узла без ручного вмешательства
Гарантия consistency над availability в условиях сетевого раздела (CP система)
Эффективная работа с poison messages через автоматический DLQ
Лучшая производительность при сетевых задержках по сравнению с mirrored queues
#Java #middle #RabbitMQ
Quorum Queues: консенсус как основа надежности
Quorum Queues используют алгоритм Raft для репликации данных между узлами кластера. Raft — это алгоритм консенсуса, обеспечивающий согласованное состояние распределенной системы при наличии отказов.
Архитектура Quorum Queue
Quorum Queue Architecture:
+-------------------+ +-------------------+ +-------------------+
| Лидер | | Последователь | | Последователь |
| (Leader) |<---->| (Follower) |<---->| (Follower) |
| | | | | |
| • Принимает запись| | • Реплицирует | | • Реплицирует |
| • Отвечает клиентам| | данные | | данные |
| • Управляет | | • Голосует за | | • Голосует за |
| логом Raft | | выборы лидера | | выборы лидера |
+-------------------+ +-------------------+ +-------------------+
| | |
| Кворум (N/2 + 1) узлов согласны |
+---------------------------------------------------+
Жизненный цикл сообщения в Quorum Queue
// Псевдокод обработки записи в Quorum Queue
quorum_queue_append(queue, message) {
// 1. Лидер добавляет запись в свой лог
log_entry = {
term: current_term,
index: next_index++,
command: "ADD_MESSAGE",
data: message
}
leader_log.append(log_entry)
// 2. Репликация на последователей
followers_acked = 1 // Лидер уже записал
for (follower in queue.followers) {
send_append_entries(follower, log_entry)
// 3. Ожидание подтверждения от большинства
if (wait_for_ack(follower, timeout)) {
followers_acked++
if (followers_acked >= quorum_size(queue)) {
// 4. Коммит записи (становится видимой для чтения)
log_entry.committed = true
apply_to_state_machine(queue, log_entry)
return SUCCESS
}
}
}
// 5. Если кворум не достигнут
if (followers_acked < quorum_size(queue)) {
// Возврат к предыдущему индексу
next_index--
return FAILURE
}
}
Преимущества Quorum Queues:
Автоматическое восстановление после потери узла без ручного вмешательства
Гарантия consistency над availability в условиях сетевого раздела (CP система)
Эффективная работа с poison messages через автоматический DLQ
Лучшая производительность при сетевых задержках по сравнению с mirrored queues
#Java #middle #RabbitMQ
👍2
Stream Queues: потоковая семантика
Stream Queues — это гибридный тип, сочетающий возможности традиционных очередей с характеристиками потоковых систем.
Ключевые особенности Stream Queues:
Append-only лог с сегментированным хранением на диске
Поддержка потребителей с разной скоростью через offset-based потребление
Дедупликация сообщений по message ID
Компрессия данных на уровне сегментов
Classic Queues: legacy с ограничениями
Classic Queues — оригинальная реализация очередей в RabbitMQ, сохраняемая для обратной совместимости:
Хранение в памяти или на диске в зависимости от persistence флагов
Mirrored queues для репликации (устаревшие, заменяются на Quorum Queues)
Ограничения при сетевых разделах (может потерять сообщения)
Более высокая производительность для ephemeral сообщений
В 2025 году использование Classic Queues рекомендуется только:
Для временных очередей (auto-delete, exclusive)
В тестовых средах
При миграции очень старых систем без возможности изменений
Стратегии хранения сообщений
Журнальная архитектура (Write-Ahead Log)
RabbitMQ использует WAL (Write-Ahead Log) для гарантированной сохранности сообщений:
Стратегии persistence в зависимости от типа очереди
Для Quorum Queues:
Все сообщения записываются на диск на всех узлах кворума
Используется segment-based хранение с индексами
Политика удержания: настраиваемое время или размер диска
Для Stream Queues:
Все сообщения хранятся на диске в сегментах
Сегменты закрываются при достижении лимита (время или размер)
Старые сегменты могут удаляться или архивироваться
Для Classic Queues с persistence:
Сообщения с флагом persistent записываются в журнал
Индекс очереди хранится в памяти для производительности
При восстановлении после сбоя: перестроение индекса из журнала
Управление памятью: ограничения и алертинг
RabbitMQ 3.13+ включает продвинутые механизмы контроля памяти:
#Java #middle #RabbitMQ
Stream Queues — это гибридный тип, сочетающий возможности традиционных очередей с характеристиками потоковых систем.
Ключевые особенности Stream Queues:
Append-only лог с сегментированным хранением на диске
Поддержка потребителей с разной скоростью через offset-based потребление
Дедупликация сообщений по message ID
Компрессия данных на уровне сегментов
// Псевдокод работы Stream Queue
stream_append(queue, message) {
// 1. Проверка дедупликации
if (queue.deduplication_enabled && message.id in seen_ids) {
return DUPLICATE
}
// 2. Добавление в текущий активный сегмент
segment = get_active_segment(queue)
// 3. Если сегмент заполнен, ротация
if (segment.size + message.size > segment.max_size) {
segment.close()
segment = create_new_segment(queue)
queue.active_segment = segment
// 4. Компрессия закрытых сегментов (асинхронно)
schedule_compression(closed_segments)
}
// 5. Запись в сегмент
write_result = segment.append(
offset: segment.next_offset++,
timestamp: current_time(),
message: message
)
// 6. Обновление индекса для быстрого поиска
update_index(queue, message.id, segment.id, write_result.position)
return write_result.offset
}
Classic Queues: legacy с ограничениями
Classic Queues — оригинальная реализация очередей в RabbitMQ, сохраняемая для обратной совместимости:
Хранение в памяти или на диске в зависимости от persistence флагов
Mirrored queues для репликации (устаревшие, заменяются на Quorum Queues)
Ограничения при сетевых разделах (может потерять сообщения)
Более высокая производительность для ephemeral сообщений
В 2025 году использование Classic Queues рекомендуется только:
Для временных очередей (auto-delete, exclusive)
В тестовых средах
При миграции очень старых систем без возможности изменений
Стратегии хранения сообщений
Журнальная архитектура (Write-Ahead Log)
RabbitMQ использует WAL (Write-Ahead Log) для гарантированной сохранности сообщений:
// Упрощенная схема работы журнала
wal_append(message) {
// 1. Запись в журнал (последовательная запись)
log_entry = serialize(message)
wal_file.append(log_entry)
fsync(wal_file) // Синхронизация с диском
// 2. Добавление в in-memory структуры для быстрого доступа
index_entry = {
message_id: message.id,
wal_position: current_position,
queue: message.queue,
status: "PERSISTED"
}
memory_index.add(index_entry)
// 3. Периодическая очистка устаревших записей
if (wal_file.size > max_wal_size) {
rotate_wal_file()
}
}
Стратегии persistence в зависимости от типа очереди
Для Quorum Queues:
Все сообщения записываются на диск на всех узлах кворума
Используется segment-based хранение с индексами
Политика удержания: настраиваемое время или размер диска
Для Stream Queues:
Все сообщения хранятся на диске в сегментах
Сегменты закрываются при достижении лимита (время или размер)
Старые сегменты могут удаляться или архивироваться
Для Classic Queues с persistence:
Сообщения с флагом persistent записываются в журнал
Индекс очереди хранится в памяти для производительности
При восстановлении после сбоя: перестроение индекса из журнала
Управление памятью: ограничения и алертинг
RabbitMQ 3.13+ включает продвинутые механизмы контроля памяти:
// Механизм flow control на основе памяти
check_memory_pressure() {
memory_used = get_memory_usage()
if (memory_used > memory_limit_high_watermark) {
// 1. Приостановка публикаций
block_publishers()
// 2. Принудительная выгрузка сообщений на диск
for (queue in queues) {
if (queue.messages_in_memory > threshold) {
page_to_disk(queue)
}
}
// 3. Если не помогло, отключение соединений
if (memory_used > memory_limit_critical) {
disconnect_clients_by_memory_usage()
}
}
}
#Java #middle #RabbitMQ
👍4
Архитектурные антипаттерны
1. Чрезмерное использование transient сообщений в production
Проблема: Использование transient (неперсистентных) сообщений для критичных данных.
Симптомы: Потеря данных при перезагрузке брокера, невоспроизводимые ошибки.
Решение: Все production очереди должны быть Quorum Queues с persistence.
2. Игнорирование flow control и back pressure
Проблема: Продьюсеры отправляют сообщения быстрее, чем консьюмеры могут их обработать.
Симптомы: Рост памяти, падение производительности, каскадные отказы.
Решение: Реализация back pressure через:
Настройку QoS prefetch count
Использование publisher confirms с таймаутами
Мониторинг глубины очередей и автоматическое масштабирование
3. Паттерн "Очередь как база данных"
Проблема: Хранение состояния приложения в очереди как primary storage.
Симптомы: Длительное восстановление после сбоя, невозможность выполнения сложных запросов.
Решение: Очередь — это транспорт, а не хранилище. Использовать специализированные БД для состояния.
4. Отсутствие стратегии dead letter обработки
Проблема: Накопление poison messages в очередях без обработки.
Симптомы: Зацикливание сообщений, блокировка обработчиков.
Решение:
5. Неправильный выбор типа очереди
Проблема: Использование Classic Queues для сценариев, требующих высокой доступности.
Симптомы: Потеря сообщений при сетевых разделах, длительные простои.
Решение:
#Java #middle #RabbitMQ
1. Чрезмерное использование transient сообщений в production
Проблема: Использование transient (неперсистентных) сообщений для критичных данных.
Симптомы: Потеря данных при перезагрузке брокера, невоспроизводимые ошибки.
Решение: Все production очереди должны быть Quorum Queues с persistence.
2. Игнорирование flow control и back pressure
Проблема: Продьюсеры отправляют сообщения быстрее, чем консьюмеры могут их обработать.
Симптомы: Рост памяти, падение производительности, каскадные отказы.
Решение: Реализация back pressure через:
Настройку QoS prefetch count
Использование publisher confirms с таймаутами
Мониторинг глубины очередей и автоматическое масштабирование
3. Паттерн "Очередь как база данных"
Проблема: Хранение состояния приложения в очереди как primary storage.
Симптомы: Длительное восстановление после сбоя, невозможность выполнения сложных запросов.
Решение: Очередь — это транспорт, а не хранилище. Использовать специализированные БД для состояния.
4. Отсутствие стратегии dead letter обработки
Проблема: Накопление poison messages в очередях без обработки.
Симптомы: Зацикливание сообщений, блокировка обработчиков.
Решение:
// Современная обработка dead letter
dead_letter_strategy {
// 1. Автоматический DLQ для Quorum Queues
policy: {
"dead-letter-strategy": "at-most-once",
"max-deliveries": 5,
"dead-letter-exchange": "dlx.global"
}
// 2. Мониторинг и алертинг DLQ
monitor_dlq_size(threshold=1000)
// 3. Автоматическая обработка или архивация
if (dlq.size > critical_threshold) {
archive_to_object_storage(dlq.messages)
dlq.purge()
}
}
5. Неправильный выбор типа очереди
Проблема: Использование Classic Queues для сценариев, требующих высокой доступности.
Симптомы: Потеря сообщений при сетевых разделах, длительные простои.
Решение:
if (требуется_гарантированная_доставка && кластер >= 3_узлов) {
использовать Quorum Queue
} else if (потоковая_обработка && хранение_истории) {
использовать Stream Queue
} else if (временные_данные && высокая_производительность) {
использовать Classic Queue in memory
} else {
// По умолчанию в 2025
использовать Quorum Queue
}#Java #middle #RabbitMQ
👍4
Что выведет код?
#Tasks
import java.util.HashSet;
public class Task291225 {
public static void main(String[] args) {
HashSet<String> set1 = new HashSet<>();
HashSet<String> set2 = new HashSet<>();
set1.add("hello");
set2.add("hello");
String s1 = new String("hello");
String s2 = new String("hello");
set1.add(s1);
set2.add(s2);
System.out.println(set1.equals(set2));
System.out.println(set1.contains(s1));
System.out.println(set2.contains(s1));
System.out.println(set1.size());
}
}
#Tasks
👍3
Варианты ответа:
Anonymous Quiz
47%
true true true 1
13%
true false true 1
20%
false true false 2
20%
true true false 2
👍2
Вопрос с собеседований
Почему утечка памяти возможна даже при наличии GC?🤓
Ответ:
GC удаляет только недостижимые объекты.
Если ссылки сохраняются, но объект логически не нужен, он останется в памяти.
Типичные причины: статические коллекции, кэши без eviction, listeners без отписки, ThreadLocal в пулах потоков.
#собеседование
Почему утечка памяти возможна даже при наличии GC?
Ответ:
Если ссылки сохраняются, но объект логически не нужен, он останется в памяти.
Типичные причины: статические коллекции, кэши без eviction, listeners без отписки, ThreadLocal в пулах потоков.
#собеседование
Please open Telegram to view this post
VIEW IN TELEGRAM
👍6
История IT-технологий сегодня — 30 декабря
ℹ️ Кто родился в этот день
Бьёрн Страуструп (устоявшееся написание; точная транскрипция дат. Bjarne Stroustrup, ˈbjɑːnə ˈsdʁʌʊ̯ˀsdʁɔb — Бьярне Строуструп; род. 30 декабря 1950, Орхус, Дания) — датский программист, автор языка программирования C++.
Кевин Систром (англ. Kevin Systrom; род. 30 декабря 1983 года, Холлистон, США) — американский предприниматель и программист, создатель и бывший генеральный директор социальной сети «Инстаграм».
🌐 Знаковые события
Не нашел(
#Biography #Birth_Date #Events #30Декабря
Бьёрн Страуструп (устоявшееся написание; точная транскрипция дат. Bjarne Stroustrup, ˈbjɑːnə ˈsdʁʌʊ̯ˀsdʁɔb — Бьярне Строуструп; род. 30 декабря 1950, Орхус, Дания) — датский программист, автор языка программирования C++.
Кевин Систром (англ. Kevin Systrom; род. 30 декабря 1983 года, Холлистон, США) — американский предприниматель и программист, создатель и бывший генеральный директор социальной сети «Инстаграм».
Не нашел(
#Biography #Birth_Date #Events #30Декабря
Please open Telegram to view this post
VIEW IN TELEGRAM
👍2🔥1
Глава 8. Дополнительные аспекты коллекций
Потокобезопасные коллекции и типичные ошибки
Многопоточное программирование представляет собой одну из наиболее сложных и тонких областей разработки программного обеспечения, где коллекции играют критически важную роль. Взаимодействие потоков через общие структуры данных требует не только технических решений, но и глубокого понимания принципов параллелизма, memory model и паттернов доступа. Потокобезопасные коллекции являются мостом между простыми однопоточными структурами данных и сложными конкурентными системами.
Эволюция подходов к потокобезопасности в Java
Исторически Java прошла несколько этапов в развитии многопоточных коллекций:
Java 1.0-1.1: Примитивная синхронизация через ключевое слово synchronized
Java 1.2: Введение Collections.synchronizedXXX() методов
Java 5 (J2SE 5.0): Революция с пакетом java.util.concurrent
Java 7-8: Усовершенствование ConcurrentHashMap и других структур
Java 9+: Дальнейшие оптимизации и новые методы
Каждый этап отражал растущее понимание сложностей многопоточного программирования и поиск баланса между производительностью, простотой использования и корректностью.
Collections.synchronizedList: Классический подход с явной синхронизацией
Collections.synchronizedList() представляет собой декоратор (wrapper) паттерн, применяемый к существующему списку для добавления потокобезопасности. Это подход минимального вмешательства — вместо создания новой потокобезопасной реализации с нуля, мы оборачиваем существующую реализацию в слой синхронизации.
Архитектура реализации
Механизм обертки
Выбор объекта монитора
Ключевое решение в дизайне — выбор объекта для синхронизации:
По умолчанию: сама обертка (this)
Альтернатива: можно передать внешний объект через конструктор SynchronizedList(list, mutex)
Это позволяет нескольким коллекциям синхронизироваться на одном мониторе, обеспечивая атомарность составных операций.
Семантика синхронизации
Уровень синхронизации
Каждый метод обертки синхронизирован индивидуально. Это обеспечивает:
Атомарность отдельных операций: Один поток не может вмешаться в выполнение метода другим потоком
Консистентность данных: Внутреннее состояние коллекции защищено от одновременных модификаций
Ограничения атомарности
Важное ограничение: хотя каждая операция атомарна, последовательность операций — нет:
Производительность и contention
Гранулярность блокировок
Collections.synchronizedList использует coarse-grained locking (грубозернистую блокировку):
Одна блокировка на всю коллекцию
Все потоки конкурируют за одну блокировку
Высокий contention при высокой конкуренции
#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList
Потокобезопасные коллекции и типичные ошибки
Многопоточное программирование представляет собой одну из наиболее сложных и тонких областей разработки программного обеспечения, где коллекции играют критически важную роль. Взаимодействие потоков через общие структуры данных требует не только технических решений, но и глубокого понимания принципов параллелизма, memory model и паттернов доступа. Потокобезопасные коллекции являются мостом между простыми однопоточными структурами данных и сложными конкурентными системами.
Эволюция подходов к потокобезопасности в Java
Исторически Java прошла несколько этапов в развитии многопоточных коллекций:
Java 1.0-1.1: Примитивная синхронизация через ключевое слово synchronized
Java 1.2: Введение Collections.synchronizedXXX() методов
Java 5 (J2SE 5.0): Революция с пакетом java.util.concurrent
Java 7-8: Усовершенствование ConcurrentHashMap и других структур
Java 9+: Дальнейшие оптимизации и новые методы
Каждый этап отражал растущее понимание сложностей многопоточного программирования и поиск баланса между производительностью, простотой использования и корректностью.
Collections.synchronizedList: Классический подход с явной синхронизацией
Collections.synchronizedList() представляет собой декоратор (wrapper) паттерн, применяемый к существующему списку для добавления потокобезопасности. Это подход минимального вмешательства — вместо создания новой потокобезопасной реализации с нуля, мы оборачиваем существующую реализацию в слой синхронизации.
Архитектура реализации
Механизм обертки
// Упрощенная концептуальная реализация
public static <T> List<T> synchronizedList(List<T> list) {
return (list instanceof RandomAccess ?
new SynchronizedRandomAccessList<>(list) :
new SynchronizedList<>(list));
}
static class SynchronizedList<E> implements List<E> {
final List<E> list; // Оборачиваемый список
final Object mutex; // Объект для синхронизации
SynchronizedList(List<E> list) {
this.list = list;
this.mutex = this; // По умолчанию синхронизируемся на обертке
}
public E get(int index) {
synchronized (mutex) { return list.get(index); }
}
public void add(int index, E element) {
synchronized (mutex) { list.add(index, element); }
}
// Все методы синхронизированы аналогично
}
Выбор объекта монитора
Ключевое решение в дизайне — выбор объекта для синхронизации:
По умолчанию: сама обертка (this)
Альтернатива: можно передать внешний объект через конструктор SynchronizedList(list, mutex)
Это позволяет нескольким коллекциям синхронизироваться на одном мониторе, обеспечивая атомарность составных операций.
Семантика синхронизации
Уровень синхронизации
Каждый метод обертки синхронизирован индивидуально. Это обеспечивает:
Атомарность отдельных операций: Один поток не может вмешаться в выполнение метода другим потоком
Консистентность данных: Внутреннее состояние коллекции защищено от одновременных модификаций
Ограничения атомарности
Важное ограничение: хотя каждая операция атомарна, последовательность операций — нет:
// ОПАСНО: неатомарная составная операция
List<String> syncList = Collections.synchronizedList(new ArrayList<>());
if (!syncList.contains("item")) { // Операция 1
syncList.add("item"); // Операция 2
}
// Между проверкой и добавлением другой поток может добавить элемент
Производительность и contention
Гранулярность блокировок
Collections.synchronizedList использует coarse-grained locking (грубозернистую блокировку):
Одна блокировка на всю коллекцию
Все потоки конкурируют за одну блокировку
Высокий contention при высокой конкуренции
#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList
👍3
Итерация и fail-fast семантика
Синхронизированные обертки не решают проблему итерации:
Для безопасной итерации требуется внешняя синхронизация:
ConcurrentHashMap: Современный подход к параллельным отображениям
ConcurrentHashMap представляет собой фундаментально иную философию по сравнению с синхронизированными обертками.
Вместо блокировки всей структуры используется комбинация:
Fine-grained locking (тонкозернистые блокировки)
Lock-free алгоритмы для чтения
CAS операции (Compare-And-Swap)
Сегментирование (в версиях до Java 8)
В Java 8 архитектура была полностью переработана:
Инновации:
CAS для вставки: sun.misc.Unsafe.compareAndSwapObject
Tree bins: Преобразование в красно-черные деревья при длинных цепочках
Параллельные операции: forEach, search, reduce
Memory model и happens-before
ConcurrentHashMap обеспечивает строгие гарантии memory ordering:
Atomicity guarantees
Параллельные операции bulk
Параметризация параллелизма
#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList
Синхронизированные обертки не решают проблему итерации:
List<String> syncList = Collections.synchronizedList(new ArrayList<>());
// ОПАСНО: ConcurrentModificationException все еще возможен
for (String item : syncList) {
// Другой поток может модифицировать список
syncList.remove("someItem"); // Из другого потока
}
Для безопасной итерации требуется внешняя синхронизация:
List<String> syncList = Collections.synchronizedList(new ArrayList<>());
// Безопасная итерация
synchronized (syncList) {
Iterator<String> it = syncList.iterator();
while (it.hasNext()) {
String item = it.next();
process(item);
}
}
ConcurrentHashMap: Современный подход к параллельным отображениям
ConcurrentHashMap представляет собой фундаментально иную философию по сравнению с синхронизированными обертками.
Вместо блокировки всей структуры используется комбинация:
Fine-grained locking (тонкозернистые блокировки)
Lock-free алгоритмы для чтения
CAS операции (Compare-And-Swap)
Сегментирование (в версиях до Java 8)
В Java 8 архитектура была полностью переработана:
// Концептуальная структура с Java 8
public class ConcurrentHashMap<K,V> {
volatile Node<K,V>[] table;
static class Node<K,V> implements Map.Entry<K,V> {
final int hash;
final K key;
volatile V val;
volatile Node<K,V> next;
}
static final class TreeNode<K,V> extends Node<K,V> {
TreeNode<K,V> parent;
TreeNode<K,V> left;
TreeNode<K,V> right;
TreeNode<K,V> prev;
boolean red;
}
}
Инновации:
CAS для вставки: sun.misc.Unsafe.compareAndSwapObject
Tree bins: Преобразование в красно-черные деревья при длинных цепочках
Параллельные операции: forEach, search, reduce
Memory model и happens-before
ConcurrentHashMap обеспечивает строгие гарантии memory ordering:
ConcurrentHashMap<String, Integer> map = new ConcurrentHashMap<>();
// Поток 1
map.put("key", 42); // Запись с memory barrier
// Путок 2
Integer value = map.get("key"); // Чтение с happens-before гарантиями
// Гарантированно увидит 42, если нет перезаписи
Atomicity guarantees
// Атомарные операции
map.putIfAbsent(key, value); // Вставить если отсутствует
map.replace(key, oldValue, newValue); // Заменить если совпадает
map.compute(key, (k, v) -> v == null ? 1 : v + 1); // Атомарное вычисление
Параллельные операции bulk
ConcurrentHashMap<String, Long> wordCounts = new ConcurrentHashMap<>();
// Параллельный forEach
wordCounts.forEach(1, // Параллелизм
(key, value) -> System.out.println(key + ":" + value));
// Поиск
String result = wordCounts.search(1,
(key, value) -> value > 1000 ? key : null);
// Свертка
long total = wordCounts.reduceValues(1, Long::sum);
Параметризация параллелизма
ConcurrentHashMap<String, Data> map = new ConcurrentHashMap<>(
16, // initial capacity
0.75f, // load factor
8 // concurrency level (оценочное количество потоков)
);
#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList
👍4
Производительность в различных сценариях
Для сценариев с частым чтением ConcurrentHashMap показывает исключительную производительность:
Чтение полностью lock-free
Минимальный contention между читателями
Эффективное использование кэшей процессора
При частой записи производительность зависит от:
Качества хэш-функции
Количества коллизий
Наличия tree bins
Конкуренции за конкретные бакеты
В конкурентной среде точный размер постоянно меняется. ConcurrentHashMap использует приближенные методы:
CopyOnWriteArrayList: Оптимизация для read-mostly сценариев
CopyOnWriteArrayList основан на фундаментальном компромиссе: дорогая запись в обмен на безопасное и эффективное чтение. Этот подход заимствован из систем управления памятью и файловых систем, где копирование при записи является стандартным паттерном.
Архитектурные принципы
Неизменяемое состояние
Ключевые особенности:
Массив объявлен как volatile для обеспечения memory visibility
Все операции чтения работают с текущим массивом
Операции записи создают новую копию
Гарантии consistency
Итераторы обеспечивают strong consistency для snapshot:
Видят состояние на момент создания
Никогда не выбрасывают ConcurrentModificationException
Не поддерживают операцию remove() (UnsupportedOperationException)
Практические паттерны использования
Event listeners и наблюдатели
Кэширование конфигураций
Ограничения и альтернативы
Когда не использовать CopyOnWriteArrayList
Частые модификации: Большие коллекции с частыми изменениями
Реальные требования: Когда нужны актуальные данные, а не snapshot
Ограничения памяти: Когда копирование больших массивов непозволительно
Альтернативные подходы
#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList
Для сценариев с частым чтением ConcurrentHashMap показывает исключительную производительность:
Чтение полностью lock-free
Минимальный contention между читателями
Эффективное использование кэшей процессора
При частой записи производительность зависит от:
Качества хэш-функции
Количества коллизий
Наличия tree bins
Конкуренции за конкретные бакеты
В конкурентной среде точный размер постоянно меняется. ConcurrentHashMap использует приближенные методы:
ConcurrentHashMap<String, String> map = new ConcurrentHashMap<>();
// Приближенный размер (O(1), но может быть неточным)
int approximateSize = map.size();
// Более точный (но дорогой) подсчет
int exactSize = map.mappingCount(); // Java 8+
// Проверка пустоты (эффективная)
boolean isEmpty = map.isEmpty();
CopyOnWriteArrayList: Оптимизация для read-mostly сценариев
CopyOnWriteArrayList основан на фундаментальном компромиссе: дорогая запись в обмен на безопасное и эффективное чтение. Этот подход заимствован из систем управления памятью и файловых систем, где копирование при записи является стандартным паттерном.
Архитектурные принципы
Неизменяемое состояние
public class CopyOnWriteArrayList<E> {
private transient volatile Object[] array;
final Object[] getArray() {
return array;
}
final void setArray(Object[] a) {
array = a;
}
}Ключевые особенности:
Массив объявлен как volatile для обеспечения memory visibility
Все операции чтения работают с текущим массивом
Операции записи создают новую копию
Гарантии consistency
Итераторы обеспечивают strong consistency для snapshot:
Видят состояние на момент создания
Никогда не выбрасывают ConcurrentModificationException
Не поддерживают операцию remove() (UnsupportedOperationException)
Практические паттерны использования
Event listeners и наблюдатели
public class EventPublisher {
private final CopyOnWriteArrayList<EventListener> listeners =
new CopyOnWriteArrayList<>();
public void addListener(EventListener listener) {
listeners.add(listener); // Безопасно даже во время уведомлений
}
public void publish(Event event) {
for (EventListener listener : listeners) {
// Итерация по snapshot - безопасна
listener.onEvent(event);
}
}
}Кэширование конфигураций
public class ConfigurationCache {
private volatile CopyOnWriteArrayList<Config> cache;
public ConfigurationCache() {
cache = new CopyOnWriteArrayList<>();
}
public void refresh() {
List<Config> newConfigs = loadConfigs();
// Атомарная замена всего кэша
cache = new CopyOnWriteArrayList<>(newConfigs);
}
public List<Config> getConfigs() {
return cache; // Безопасное чтение
}
}Ограничения и альтернативы
Когда не использовать CopyOnWriteArrayList
Частые модификации: Большие коллекции с частыми изменениями
Реальные требования: Когда нужны актуальные данные, а не snapshot
Ограничения памяти: Когда копирование больших массивов непозволительно
Альтернативные подходы
// Для частых модификаций
List<String> frequentWrites = Collections.synchronizedList(new ArrayList<>());
// Для mixed workloads
ConcurrentLinkedQueue<String> queue = new ConcurrentLinkedQueue<>();
// Для сценариев с преобладанием чтения
List<String> readMostly = new CopyOnWriteArrayList<>();
#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList
👍3
Типичные ошибки многопоточного программирования с коллекциями
ConcurrentModificationException: Анатомия ошибки
ConcurrentModificationException возникает при обнаружении структурных изменений коллекции во время итерации.
Механизм основан на сравнении счетчика модификаций:
Типичные сценарии возникновения
Сценарий 1: Модификация во время итерации в одном потоке
Сценарий 2: Конкурентная модификация в разных потоках
NullPointerException в многопоточном контексте
NullPointerException в многопоточных сценариях часто является следствием race conditions, а не просто нулевых ссылок:
Классический антипаттерн с небезопасной публикацией:
Разные коллекции по-разному обрабатывают null:
#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList
ConcurrentModificationException: Анатомия ошибки
ConcurrentModificationException возникает при обнаружении структурных изменений коллекции во время итерации.
Механизм основан на сравнении счетчика модификаций:
// Внутренний механизм ArrayList
protected transient int modCount = 0;
// В итераторе
int expectedModCount = modCount;
void checkForComodification() {
if (modCount != expectedModCount)
throw new ConcurrentModificationException();
}
Типичные сценарии возникновения
Сценарий 1: Модификация во время итерации в одном потоке
List<String> list = new ArrayList<>(Arrays.asList("A", "B", "C"));
for (String item : list) { // Создается итератор
if (item.equals("B")) {
list.remove(item); // modCount++ → исключение!
}
}Сценарий 2: Конкурентная модификация в разных потоках
// Поток 1
for (String item : sharedList) {
process(item); // Итерация
}
// Поток 2
sharedList.add("new"); // ConcurrentModificationException в потоке 1
NullPointerException в многопоточном контексте
NullPointerException в многопоточных сценариях часто является следствием race conditions, а не просто нулевых ссылок:
public class UnsafeCache {
private Map<String, Data> cache = new HashMap<>();
public Data get(String key) {
Data data = cache.get(key);
if (data == null) {
data = loadData(key); // Дорогая операция
cache.put(key, data); // Race condition!
}
return data; // Может вернуть null
}
}Классический антипаттерн с небезопасной публикацией:
public class BrokenSingleton {
private static Data instance;
public static Data getInstance() {
if (instance == null) { // Первая проверка (без синхронизации)
synchronized (BrokenSingleton.class) {
if (instance == null) { // Вторая проверка
instance = new Data(); // Небезопасная публикация!
}
}
}
return instance; // Может вернуть частично инициализированный объект
}
}Разные коллекции по-разному обрабатывают null:
// ConcurrentHashMap: запрещает null
ConcurrentHashMap<String, String> chm = new ConcurrentHashMap<>();
chm.put("key", null); // NullPointerException
// CopyOnWriteArrayList: разрешает null
CopyOnWriteArrayList<String> cowal = new CopyOnWriteArrayList<>();
cowal.add(null); // Допустимо
// Collections.synchronizedList: зависит от оборачиваемой коллекции
List<String> syncList = Collections.synchronizedList(new ArrayList<>());
syncList.add(null); // Допустимо для ArrayList
#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList
👍2
Race conditions и data races
Race condition: Неправильное поведение из-за непредсказуемого порядка выполнения
Data race: Одновременный доступ к shared memory без proper synchronization
Пример race condition
Пример data race
Deadlock, livelock и starvation
Deadlock с коллекциями
Livelock в конкурентных алгоритмах
Starvation в synchronized коллекциях
#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList
Race condition: Неправильное поведение из-за непредсказуемого порядка выполнения
Data race: Одновременный доступ к shared memory без proper synchronization
Пример race condition
public class Counter {
private int count;
public void increment() {
count++; // Неатомарная операция: read-modify-write
}
}
// Два потока вызывают increment() 1000 раз каждый
// Ожидаемый результат: 2000
// Фактический результат: что угодно между 1000 и 2000Пример data race
public class VisibilityProblem {
private boolean ready = false;
private int value;
// Поток 1
public void writer() {
value = 42;
ready = true; // Без happens-before!
}
// Путок 2
public void reader() {
if (ready) {
System.out.println(value); // Может увидеть 0 вместо 42!
}
}
}Deadlock, livelock и starvation
Deadlock с коллекциями
// Классический deadlock с synchronizedList
List<String> list1 = Collections.synchronizedList(new ArrayList<>());
List<String> list2 = Collections.synchronizedList(new ArrayList<>());
// Поток 1
synchronized (list1) {
synchronized (list2) { // Ждет list2
// Критическая секция
}
}
// Путок 2 (обратный порядок)
synchronized (list2) {
synchronized (list1) { // Ждет list1 → DEADLOCK!
// Критическая секция
}
}
Livelock в конкурентных алгоритмах
public class LivelockExample {
private final ConcurrentHashMap<String, Boolean> locks =
new ConcurrentHashMap<>();
public void process(String key) {
// Бесконечные попытки захвата "локера"
while (!locks.putIfAbsent(key, true)) {
Thread.yield(); // Livelock: постоянно уступаем, но не прогрессируем
}
try {
// Работа с ресурсом
} finally {
locks.remove(key);
}
}
}Starvation в synchronized коллекциях
// Поток, постоянно читающий
synchronized (sharedList) {
// Долгая операция чтения
processAllElements(sharedList);
}
// Другие потоки не могут получить доступ для записи
// → Starvation писателей
#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList
👍2
Что выведет код?
#Tasks
import java.util.concurrent.*;
public class Task301225 {
public static void main(String[] args) {
ConcurrentHashMap<String, Integer> map = new ConcurrentHashMap<>();
map.put("a", 1);
Integer result = map.computeIfAbsent("b", k -> {
map.put("b", 2);
return map.get("b");
});
System.out.println(result);
System.out.println(map.get("b"));
}
}
#Tasks
👍1
👍1
