Description:
mysql::csa::Cached_event_payload
(sql/changestreams/apply/storage/relay_log/cached_event_payload.{h,cpp})
owns a raw event buffer, uint8_t *m_data, allocated by
Default_binlog_event_allocator (instrumented as memory/sql/Log_event).
The buffer is released only inside decode():
- on success, it is handed to the returned Log_event
(register_temp_buf with DELEGATE_MEMORY_TO_EVENT_OBJECT);
- on error, it is freed with m_allocator.deallocate().
The class has no destructor. If a Cached_event_payload is destroyed
before decode() is called, m_data is never freed and the buffer leaks.
The CSA applier stores Cached_event_payload objects in
Event_set_fetchable_cache batches (cache_metadata read mode). Some of
these objects are destroyed without being decoded, for example:
- Job_applier::run_phase() finds the transaction's GTID already in
gtid_executed, calls set_fetching_done(), and never fetches the
remaining events of the transaction.
- A queued job is discarded on STOP REPLICA or when the reader
truncates a transaction.
- Event_set_fetchable_cache::append_event() drops the event because the
batch is already sealed or truncated.
Every undecoded event in these cases leaks one buffer of the event's
size.
The default build does not show the leak:
Relay_log_adaptive_reader::read() uses cache_metadata mode only when
trx_size > provider_max_read_event_bytes &&
trx_size < provider_max_read_payload_bytes
Both constants in sql/changestreams/apply/context/tune.h are 512, and
the runtime auto-tuning only switches them between 512/512 and 0/0.
The range is always empty, so no Cached_event_payload is created.
How to repeat:
0. Build mysqld with this change so the reader uses cache_metadata mode:
```
--- a/sql/changestreams/apply/context/tune.h
+++ b/sql/changestreams/apply/context/tune.h
@@ -44,8 +44,8 @@
inline constexpr std::size_t provider_max_read_event_bytes{512};
-inline constexpr std::size_t provider_max_read_payload_bytes{512};
-inline constexpr bool csa_provider_enable_tune{true};
+inline constexpr std::size_t provider_max_read_payload_bytes{64 * 1024 * 1024};
+inline constexpr bool csa_provider_enable_tune{false};
```
The first change makes cache_metadata cover transactions from 512
bytes to 64 MiB. The second stops auto-tuning from resetting the
thresholds to 0 under low worker load.
1. Set up a source and a replica with GTID_MODE=ON and binlog_format=ROW.
On the replica:
CHANGE REPLICATION SOURCE TO APPLIER_VERSION=2, APPLIER_WORKER_COUNT=4;
START REPLICA;
On the source:
CREATE TABLE test.t1 (id INT PRIMARY KEY, b LONGBLOB);
Wait for the replica to catch up.
2. On the replica, record the baseline:
SELECT CURRENT_NUMBER_OF_BYTES_USED
FROM performance_schema.memory_summary_global_by_event_name
WHERE EVENT_NAME = 'memory/sql/Log_event';
3. Repeat 20 times with N = 1001..1020:
-- replica: pre-execute the GTID so the applier skips the transaction
SET GTID_NEXT='<source_uuid>:N'; BEGIN; COMMIT; SET GTID_NEXT=AUTOMATIC;
-- source: ~500 KB transaction with the same GTID
SET GTID_NEXT='<source_uuid>:N';
BEGIN;
INSERT INTO t1 VALUES (N*100+1, REPEAT('a',50000));
INSERT INTO t1 VALUES (N*100+2, REPEAT('a',50000));
... (10 rows in total, ids N*100+1 .. N*100+10)
COMMIT;
SET GTID_NEXT=AUTOMATIC;
4. Commit one more ordinary transaction on the source, for example
INSERT INTO t1 VALUES (1, 'marker');
and wait for the replica to catch up.
5. On the replica, run the query from step 2 again.
Result: the value grew by 10,033,880 bytes, about the total size of the
skipped row events, and does not go back down. With the suggested fix,
the growth is 0 bytes.
Suggested fix:
Add a destructor that frees m_data if decode() never handed it off, and
clear m_data before returning on the decode() error path so the
destructor cannot free it again.
```
--- a/sql/changestreams/apply/storage/relay_log/cached_event_payload.h
+++ b/sql/changestreams/apply/storage/relay_log/cached_event_payload.h
@@ class Cached_event_payload : public IReader_event {
Cached_event_payload(const Event_payload &payload,
std::shared_ptr<Log_event> fde);
+ /// @brief Frees the owned payload buffer if it was never handed off by
+ /// decode().
+ ~Cached_event_payload() override;
+
--- a/sql/changestreams/apply/storage/relay_log/cached_event_payload.cpp
+++ b/sql/changestreams/apply/storage/relay_log/cached_event_payload.cpp
+Cached_event_payload::~Cached_event_payload() {
+ if (m_data != nullptr) {
+ m_allocator.deallocate(m_data);
+ m_data = nullptr;
+ }
+}
+
@@ std::shared_ptr<Log_event> Cached_event_payload::decode() {
if (read_status.has_error()) {
m_allocator.deallocate(m_data);
- return std::shared_ptr<Log_event>();
m_data = nullptr;
+ return std::shared_ptr<Log_event>();
}
```
Description: mysql::csa::Cached_event_payload (sql/changestreams/apply/storage/relay_log/cached_event_payload.{h,cpp}) owns a raw event buffer, uint8_t *m_data, allocated by Default_binlog_event_allocator (instrumented as memory/sql/Log_event). The buffer is released only inside decode(): - on success, it is handed to the returned Log_event (register_temp_buf with DELEGATE_MEMORY_TO_EVENT_OBJECT); - on error, it is freed with m_allocator.deallocate(). The class has no destructor. If a Cached_event_payload is destroyed before decode() is called, m_data is never freed and the buffer leaks. The CSA applier stores Cached_event_payload objects in Event_set_fetchable_cache batches (cache_metadata read mode). Some of these objects are destroyed without being decoded, for example: - Job_applier::run_phase() finds the transaction's GTID already in gtid_executed, calls set_fetching_done(), and never fetches the remaining events of the transaction. - A queued job is discarded on STOP REPLICA or when the reader truncates a transaction. - Event_set_fetchable_cache::append_event() drops the event because the batch is already sealed or truncated. Every undecoded event in these cases leaks one buffer of the event's size. The default build does not show the leak: Relay_log_adaptive_reader::read() uses cache_metadata mode only when trx_size > provider_max_read_event_bytes && trx_size < provider_max_read_payload_bytes Both constants in sql/changestreams/apply/context/tune.h are 512, and the runtime auto-tuning only switches them between 512/512 and 0/0. The range is always empty, so no Cached_event_payload is created. How to repeat: 0. Build mysqld with this change so the reader uses cache_metadata mode: ``` --- a/sql/changestreams/apply/context/tune.h +++ b/sql/changestreams/apply/context/tune.h @@ -44,8 +44,8 @@ inline constexpr std::size_t provider_max_read_event_bytes{512}; -inline constexpr std::size_t provider_max_read_payload_bytes{512}; -inline constexpr bool csa_provider_enable_tune{true}; +inline constexpr std::size_t provider_max_read_payload_bytes{64 * 1024 * 1024}; +inline constexpr bool csa_provider_enable_tune{false}; ``` The first change makes cache_metadata cover transactions from 512 bytes to 64 MiB. The second stops auto-tuning from resetting the thresholds to 0 under low worker load. 1. Set up a source and a replica with GTID_MODE=ON and binlog_format=ROW. On the replica: CHANGE REPLICATION SOURCE TO APPLIER_VERSION=2, APPLIER_WORKER_COUNT=4; START REPLICA; On the source: CREATE TABLE test.t1 (id INT PRIMARY KEY, b LONGBLOB); Wait for the replica to catch up. 2. On the replica, record the baseline: SELECT CURRENT_NUMBER_OF_BYTES_USED FROM performance_schema.memory_summary_global_by_event_name WHERE EVENT_NAME = 'memory/sql/Log_event'; 3. Repeat 20 times with N = 1001..1020: -- replica: pre-execute the GTID so the applier skips the transaction SET GTID_NEXT='<source_uuid>:N'; BEGIN; COMMIT; SET GTID_NEXT=AUTOMATIC; -- source: ~500 KB transaction with the same GTID SET GTID_NEXT='<source_uuid>:N'; BEGIN; INSERT INTO t1 VALUES (N*100+1, REPEAT('a',50000)); INSERT INTO t1 VALUES (N*100+2, REPEAT('a',50000)); ... (10 rows in total, ids N*100+1 .. N*100+10) COMMIT; SET GTID_NEXT=AUTOMATIC; 4. Commit one more ordinary transaction on the source, for example INSERT INTO t1 VALUES (1, 'marker'); and wait for the replica to catch up. 5. On the replica, run the query from step 2 again. Result: the value grew by 10,033,880 bytes, about the total size of the skipped row events, and does not go back down. With the suggested fix, the growth is 0 bytes. Suggested fix: Add a destructor that frees m_data if decode() never handed it off, and clear m_data before returning on the decode() error path so the destructor cannot free it again. ``` --- a/sql/changestreams/apply/storage/relay_log/cached_event_payload.h +++ b/sql/changestreams/apply/storage/relay_log/cached_event_payload.h @@ class Cached_event_payload : public IReader_event { Cached_event_payload(const Event_payload &payload, std::shared_ptr<Log_event> fde); + /// @brief Frees the owned payload buffer if it was never handed off by + /// decode(). + ~Cached_event_payload() override; + --- a/sql/changestreams/apply/storage/relay_log/cached_event_payload.cpp +++ b/sql/changestreams/apply/storage/relay_log/cached_event_payload.cpp +Cached_event_payload::~Cached_event_payload() { + if (m_data != nullptr) { + m_allocator.deallocate(m_data); + m_data = nullptr; + } +} + @@ std::shared_ptr<Log_event> Cached_event_payload::decode() { if (read_status.has_error()) { m_allocator.deallocate(m_data); - return std::shared_ptr<Log_event>(); m_data = nullptr; + return std::shared_ptr<Log_event>(); } ```