diff --git a/snuba/cli/bulk_load.py b/snuba/cli/bulk_load.py index 8525b6e7d96..95f45f9aa70 100644 --- a/snuba/cli/bulk_load.py +++ b/snuba/cli/bulk_load.py @@ -5,6 +5,7 @@ from snuba import settings from snuba.datasets.factory import get_dataset, DATASET_NAMES from snuba.snapshots.postgres_snapshot import PostgresSnapshot +from snuba.writer import BufferedWriterWrapper @click.command() @@ -32,5 +33,9 @@ def bulk_load(dataset, dest_table, source, log_level): ) loader = dataset.get_bulk_loader(snapshot_source, dest_table) - - loader.load() + # TODO: see whether we need to pass options to the writer + writer = BufferedWriterWrapper( + dataset.get_writer({}, dest_table), + settings.BULK_CLICKHOUSE_BUFFER, + ) + loader.load(writer) diff --git a/snuba/datasets/__init__.py b/snuba/datasets/__init__.py index e3c19230349..8fe65864b5a 100644 --- a/snuba/datasets/__init__.py +++ b/snuba/datasets/__init__.py @@ -27,7 +27,7 @@ def get_schema(self): def get_processor(self): return self.__processor - def get_writer(self, options=None): + def get_writer(self, options=None, table_name=None): from snuba import settings from snuba.writer import HTTPBatchWriter @@ -36,6 +36,7 @@ def get_writer(self, options=None): settings.CLICKHOUSE_HOST, settings.CLICKHOUSE_HTTP_PORT, options, + table_name, ) def default_conditions(self): diff --git a/snuba/datasets/cdc/cdcprocessors.py b/snuba/datasets/cdc/cdcprocessors.py index 004501e712d..76476ea2452 100644 --- a/snuba/datasets/cdc/cdcprocessors.py +++ b/snuba/datasets/cdc/cdcprocessors.py @@ -1,26 +1,92 @@ +from __future__ import annotations + +from abc import ABC, abstractmethod +from typing import Any, Mapping, Optional, Sequence, Type + from snuba.processor import MessageProcessor +from snuba.writer import WriterTableRow KAFKA_ONLY_PARTITION = 0 # CDC only works with single partition topics. So partition must be 0 +class CdcMessageRow(ABC): + """ + Takes care of the data transformation from WAL to clickhouse and from + bulk load to clickhouse. The goal is to keep all these transformation + function in the same place because they ultimately have to be consistent + with the Clickhouse schema. + """ + + @classmethod + def from_wal( + cls, + offset: int, + columnnames: Sequence[str], + columnvalues: Sequence[Any], + ) -> CdcMessageRow: + raise NotImplementedError + + @classmethod + def from_bulk( + cls, + row: Mapping[str, Any], + ) -> CdcMessageRow: + raise NotImplementedError + + @abstractmethod + def to_clickhouse(self) -> WriterTableRow: + raise NotImplementedError + + class CdcProcessor(MessageProcessor): - def __init__(self, pg_table): + def __init__(self, pg_table: str, message_row_class: Type[CdcMessageRow]): self.pg_table = pg_table + self._message_row_class = message_row_class - def _process_begin(self, offset): + def _process_begin(self, offset: int): pass - def _process_commit(self, offset): + def _process_commit(self, offset: int): pass - def _process_insert(self, offset, columnnames, columnvalues): - pass + def _process_insert(self, + offset: int, + columnnames: Sequence[str], + columnvalues: Sequence[Any], + ) -> Optional[WriterTableRow]: + return self._message_row_class.from_wal( + offset, + columnnames, + columnvalues + ).to_clickhouse() - def _process_update(self, offset, key, columnnames, columnvalues): - pass + def _process_update(self, + offset: int, + key: Mapping[str, Any], + columnnames: Sequence[str], + columnvalues: Sequence[Any], + ) -> Optional[WriterTableRow]: + old_key = dict(zip(key['keynames'], key['keyvalues'])) + new_key = { + key: columnvalues[columnnames.index(key)] + for key + in key['keynames'] + } + # We cannot support a change in the identity of the record + # clickhouse will use the identity column to find rows to merge. + # if we change it, merging won't work. + assert old_key == new_key, 'Changing Primary Key is not supported.' + return self._message_row_class.from_wal( + offset, + columnnames, + columnvalues + ).to_clickhouse() - def _process_delete(self, offset, key): + def _process_delete(self, + offset: int, + key: Mapping[str, Any], + ) -> Optional[WriterTableRow]: pass def process_message(self, value, metadata): diff --git a/snuba/datasets/cdc/groupedmessage.py b/snuba/datasets/cdc/groupedmessage.py index a10cb1f6da2..679a4195149 100644 --- a/snuba/datasets/cdc/groupedmessage.py +++ b/snuba/datasets/cdc/groupedmessage.py @@ -1,6 +1,6 @@ from snuba.clickhouse import ColumnSet, DateTime, Nullable, UInt from snuba.datasets import Dataset -from snuba.datasets.cdc.groupedmessage_processor import GroupedMessageProcessor +from snuba.datasets.cdc.groupedmessage_processor import GroupedMessageProcessor, GroupedMessageRow from snuba.datasets.schema import ReplacingMergeTreeSchema from snuba.snapshots.bulk_load import SingleTableBulkLoader @@ -55,4 +55,5 @@ def get_bulk_loader(self, source, dest_table): source=source, source_table=self.POSTGRES_TABLE, dest_table=dest_table, + row_processor=lambda row: GroupedMessageRow.from_bulk(row).to_clickhouse(), ) diff --git a/snuba/datasets/cdc/groupedmessage_processor.py b/snuba/datasets/cdc/groupedmessage_processor.py index 7262ad04136..e1cbc4132ba 100644 --- a/snuba/datasets/cdc/groupedmessage_processor.py +++ b/snuba/datasets/cdc/groupedmessage_processor.py @@ -1,54 +1,99 @@ +from __future__ import annotations + +from datetime import datetime +from typing import Any, Mapping, Optional, Sequence + +from dataclasses import dataclass from dateutil.parser import parse as dateutil_parse -from snuba.datasets.cdc.cdcprocessors import CdcProcessor +from snuba.datasets.cdc.cdcprocessors import CdcProcessor, CdcMessageRow +from snuba.writer import WriterTableRow -class GroupedMessageProcessor(CdcProcessor): +@dataclass(frozen=True) +class GroupMessageRecord: + status: int + last_seen: datetime + first_seen: datetime + active_at: Optional[datetime] = None + first_release_id: Optional[int] = None - def __init__(self, postgres_table): - super(GroupedMessageProcessor, self).__init__( - pg_table=postgres_table, - ) - def __build_record(self, offset, columnnames, columnvalues): - raw_data = dict(zip(columnnames, columnvalues)) - output = { - 'offset': offset, - 'id': raw_data['id'], - 'record_deleted': 0, - 'status': raw_data['status'], - 'last_seen': dateutil_parse(raw_data['last_seen']), - 'first_seen': dateutil_parse(raw_data['first_seen']), - 'active_at': dateutil_parse(raw_data['active_at']), - 'first_release_id': raw_data['first_release_id'], - } +@dataclass(frozen=True) +class GroupedMessageRow(CdcMessageRow): + offset: Optional[int] + id: int + record_deleted: bool + record_content: Optional[GroupMessageRecord] - return output + @classmethod + def from_wal(cls, + offset: int, + columnnames: Sequence[str], + columnvalues: Sequence[Any], + ) -> GroupedMessageRow: + raw_data = dict(zip(columnnames, columnvalues)) + return GroupedMessageRow( + offset=offset, + id=raw_data['id'], + record_deleted=False, + record_content=GroupMessageRecord( + status=raw_data['status'], + last_seen=dateutil_parse(raw_data['last_seen']), + first_seen=dateutil_parse(raw_data['first_seen']), + active_at=dateutil_parse(raw_data['active_at']), + first_release_id=raw_data['first_release_id'], + ) + ) - def _process_insert(self, offset, columnnames, columnvalues): - return self.__build_record(offset, columnnames, columnvalues) + @classmethod + def from_bulk(cls, + row: Mapping[str, Any], + ) -> GroupedMessageRow: + return GroupedMessageRow( + offset=None, + id=int(row['id']), + record_deleted=False, + record_content=GroupMessageRecord( + status=int(row['status']), + last_seen=dateutil_parse(row['last_seen']), + first_seen=dateutil_parse(row['first_seen']), + active_at=dateutil_parse(row['active_at']), + first_release_id=int(row['first_release_id']) if row['first_release_id'] else None, + ) + ) - def _process_delete(self, offset, key): - key_names = key['keynames'] - key_values = key['keyvalues'] - id = key_values[key_names.index('id')] + def to_clickhouse(self) -> WriterTableRow: + record = self.record_content return { - 'offset': offset, - 'id': id, - 'record_deleted': 1, - 'status': None, - 'last_seen': None, - 'first_seen': None, - 'active_at': None, - 'first_release_id': None, + 'offset': self.offset if self.offset is not None else 0, + 'id': self.id, + 'record_deleted': 1 if self.record_deleted else 0, + 'status': None if not record else record.status, + 'last_seen': None if not record else record.last_seen, + 'first_seen': None if not record else record.first_seen, + 'active_at': None if not record else record.active_at, + 'first_release_id': None if not record else record.first_release_id, } - def _process_update(self, offset, key, columnnames, columnvalues): - new_id = columnvalues[columnnames.index('id')] + +class GroupedMessageProcessor(CdcProcessor): + + def __init__(self, postgres_table): + super(GroupedMessageProcessor, self).__init__( + pg_table=postgres_table, + message_row_class=GroupedMessageRow, + ) + + def _process_delete(self, + offset: int, + key: Mapping[str, Any], + ) -> Optional[WriterTableRow]: key_names = key['keynames'] key_values = key['keyvalues'] - old_id = key_values[key_names.index('id')] - # We cannot support a change in the identity of the record - # clickhouse will use the identity column to find rows to merge. - # if we change it, merging won't work. - assert old_id == new_id, 'Changing Primary Key is not supported.' - return self.__build_record(offset, columnnames, columnvalues) + id = key_values[key_names.index('id')] + return GroupedMessageRow( + offset=offset, + id=id, + record_deleted=True, + record_content=None + ).to_clickhouse() diff --git a/snuba/settings_base.py b/snuba/settings_base.py index ec4c2c5c3ea..7d96fa15065 100644 --- a/snuba/settings_base.py +++ b/snuba/settings_base.py @@ -43,6 +43,7 @@ # Snuba Options SNAPSHOT_LOAD_PRODUCT = 'snuba' +BULK_CLICKHOUSE_BUFFER = 1000 # Processor/Writer Options DEFAULT_BROKERS = ['localhost:9092'] diff --git a/snuba/snapshots/__init__.py b/snuba/snapshots/__init__.py index af8a761672e..9cd5b943a5e 100644 --- a/snuba/snapshots/__init__.py +++ b/snuba/snapshots/__init__.py @@ -1,12 +1,14 @@ from __future__ import annotations from abc import ABC, abstractmethod +from collections.abc import Iterator from contextlib import contextmanager -from typing import NewType, Generator, IO, Optional, Sequence +from typing import Any, Mapping, NewType, Generator, IO, Iterable, Optional, Sequence from dataclasses import dataclass SnapshotId = NewType("SnapshotId", str) +SnapshotTableRow = Mapping[str, Any] @dataclass(frozen=True) @@ -26,6 +28,12 @@ class SnapshotDescriptor: id: SnapshotId tables: Sequence[TableConfig] + def get_table(self, table_name: str): + for t in self.tables: + if t.table == table_name: + return t + raise ValueError(f"Table {table_name} does not exists in the snapshot") + class BulkLoadSource(ABC): """ @@ -41,5 +49,5 @@ def get_descriptor(self) -> SnapshotDescriptor: @abstractmethod @contextmanager - def get_table_file(self, table: str) -> Generator[IO[bytes], None, None]: + def get_table_file(self, table: str) -> Generator[Iterable[SnapshotTableRow], None, None]: raise NotImplementedError diff --git a/snuba/snapshots/bulk_load.py b/snuba/snapshots/bulk_load.py index 41ef578ec26..d94bf57698b 100644 --- a/snuba/snapshots/bulk_load.py +++ b/snuba/snapshots/bulk_load.py @@ -1,8 +1,12 @@ from abc import ABC, abstractmethod +from typing import Any, Callable, Mapping + import logging from snuba.clickhouse import ClickhousePool from snuba.snapshots import BulkLoadSource +from snuba.writer import BufferedWriterWrapper, WriterTableRow +from snuba.snapshots import SnapshotTableRow class BulkLoader(ABC): @@ -14,7 +18,7 @@ class BulkLoader(ABC): the bulk load operation. """ @abstractmethod - def load(self) -> None: + def load(self, writer: BufferedWriterWrapper) -> None: raise NotImplementedError @@ -24,15 +28,17 @@ class SingleTableBulkLoader(BulkLoader): """ def __init__(self, - source: BulkLoadSource, - dest_table: str, - dataset_table: str, - ): + source: BulkLoadSource, + dest_table: str, + source_table: str, + row_processor: Callable[[SnapshotTableRow], WriterTableRow], + ): self.__source = source self.__dest_table = dest_table self.__source_table = source_table + self.__row_processor = row_processor - def load(self) -> None: + def load(self, writer: BufferedWriterWrapper) -> None: logger = logging.getLogger('snuba.bulk-loader') clickhouse_ro = ClickhousePool(client_settings={ @@ -50,6 +56,11 @@ def load(self) -> None: logger.info("Loading snapshot %s", descriptor.id) with self.__source.get_table_file(self.__source_table) as table: - logger.info("Loading table from file %s", table.name) - # TODO: Do something with the table file - raise NotImplementedError + logger.info("Loading table %s from file", self.__source_table) + row_count = 0 + with writer as buffer_writer: + for row in table: + clickhouse_data = self.__row_processor(row) + buffer_writer.write(clickhouse_data) + row_count += 1 + logger.info("Load complete %d records loaded", row_count) diff --git a/snuba/snapshots/postgres_snapshot.py b/snuba/snapshots/postgres_snapshot.py index 08ec4566c79..71688ccf889 100644 --- a/snuba/snapshots/postgres_snapshot.py +++ b/snuba/snapshots/postgres_snapshot.py @@ -1,5 +1,6 @@ from __future__ import annotations +import csv import jsonschema # type: ignore import json import logging @@ -7,10 +8,10 @@ from contextlib import contextmanager from dataclasses import dataclass -from typing import NewType, Generator, IO, Sequence +from typing import Any, Mapping, NewType, Generator, Iterable, Sequence from snuba.snapshots import SnapshotDescriptor, TableConfig -from snuba.snapshots import BulkLoadSource +from snuba.snapshots import BulkLoadSource, SnapshotTableRow Xid = NewType("Xid", int) @@ -121,11 +122,39 @@ def get_descriptor(self) -> PostgresSnapshotDescriptor: return self.__descriptor @contextmanager - def get_table_file(self, table: str) -> Generator[IO[bytes], None, None]: + def get_table_file( + self, + table: str, + ) -> Generator[Iterable[SnapshotTableRow], None, None]: table_path = os.path.join(self.__path, "tables", "%s.csv" % table) try: - with open(table_path, "rb") as table_file: - yield table_file + with open(table_path, "r") as table_file: + csv_file = csv.DictReader(table_file) + columns = csv_file.fieldnames + + expected_columns = self.__descriptor.get_table(table).columns + if expected_columns: + expected_set = set(expected_columns) + existing_set = set(columns) + if not expected_set <= existing_set: + raise ValueError( + "The table %s is missing columns %r " % ( + table, + expected_set - existing_set, + ) + ) + + if len(existing_set) != len(expected_set): + logger.warning( + "The table %s contains more columns than expected %r", + table, + existing_set - expected_set, + ) + else: + logger.info("Won't pre-validate snapshot columns. There is nothing in the descriptor") + + yield csv_file + except FileNotFoundError: raise ValueError( "The snapshot does not contain the requested table %s" % table, diff --git a/snuba/writer.py b/snuba/writer.py index e3d44bb387d..f0308b4fd92 100644 --- a/snuba/writer.py +++ b/snuba/writer.py @@ -1,6 +1,7 @@ import json import logging from datetime import datetime +from typing import Any, Iterable, List, Mapping from urllib3.connectionpool import HTTPConnectionPool from urllib3.exceptions import HTTPError @@ -11,6 +12,8 @@ logger = logging.getLogger('snuba.writer') +WriterTableRow = Mapping[str, Any] + class BatchWriter(object): def __init__(self, schema): @@ -34,7 +37,7 @@ def __row_to_column_list(self, columns, row): values.append(value) return values - def write(self, rows): + def write(self, rows: Iterable[WriterTableRow]): columns = self.__schema.get_columns() self.__connection.execute_robust("INSERT INTO %(table)s (%(colnames)s) VALUES" % { 'colnames': ", ".join(col.escaped for col in columns), @@ -43,10 +46,11 @@ def write(self, rows): class HTTPBatchWriter(BatchWriter): - def __init__(self, schema, host, port, options=None): + def __init__(self, schema, host, port, options=None, table_name=None): self.__schema = schema self.__pool = HTTPConnectionPool(host, port) self.__options = options if options is not None else {} + self.__table_name = table_name or schema.get_table_name() def __default(self, value): if isinstance(value, datetime): @@ -54,12 +58,12 @@ def __default(self, value): else: raise TypeError - def __encode(self, row): + def __encode(self, row: WriterTableRow): return json.dumps(row, default=self.__default).encode('utf-8') - def write(self, rows): + def write(self, rows: Iterable[WriterTableRow]): parameters = self.__options.copy() - parameters['query'] = f"INSERT INTO {self.__schema.get_table_name()} FORMAT JSONEachRow" + parameters['query'] = f"INSERT INTO {self.__table_name} FORMAT JSONEachRow" resp = self.__pool.urlopen( 'POST', '/?' + urlencode(parameters), @@ -70,5 +74,38 @@ def write(self, rows): body=map(self.__encode, rows), chunked=True, ) - if resp.status//100 != 2: + if resp.status // 100 != 2: raise HTTPError(f"{resp.status} Unexpected") + + +class BufferedWriterWrapper: + """ + This is a wrapper that adds a buffer around a BatchWriter. + When consuming data from Kafka, the buffering logic is performed by the + batching consumer. + This is for the use cases that are not Kafka related. + + This is not thread safe. Don't try to do parallel flush hoping in the GIL. + """ + + def __init__(self, writer: BatchWriter, buffer_size: int): + self.__writer = writer + self.__buffer_size = buffer_size + self.__buffer: List[WriterTableRow] = [] + + def __flush(self) -> None: + logger.debug("Flushing buffer with %d elements", len(self.__buffer)) + self.__writer.write(self.__buffer) + self.__buffer = [] + + def __enter__(self): + return self + + def __exit__(self, type, value, traceback): + if self.__buffer: + self.__flush() + + def write(self, row: WriterTableRow): + self.__buffer.append(row) + if len(self.__buffer) >= self.__buffer_size: + self.__flush() diff --git a/tests/base.py b/tests/base.py index d24b9d3a2d2..f3b4d5f3cc5 100644 --- a/tests/base.py +++ b/tests/base.py @@ -30,6 +30,7 @@ def wrap_raw_event(event): 'data': event } + def get_event(): from fixtures import raw_event timestamp = datetime.utcnow() @@ -141,7 +142,24 @@ def teardown_method(self, test_method): redis_client.flushdb() -class BaseEventsTest(BaseTest): +class BaseDatasetTest(BaseTest): + def write_processed_records(self, records): + if not isinstance(records, (list, tuple)): + records = [records] + + rows = [] + for event in records: + rows.append(event) + + return self.write_rows(rows) + + def write_rows(self, rows): + if not isinstance(rows, (list, tuple)): + rows = [rows] + self.dataset.get_writer().write(rows) + + +class BaseEventsTest(BaseDatasetTest): def setup_method(self, test_method): super(BaseEventsTest, self).setup_method(test_method, 'events') self.event = get_event() @@ -168,20 +186,4 @@ def write_raw_events(self, events): _, processed = self.dataset.get_processor().process_message(event) out.append(processed) - return self.write_processed_events(out) - - def write_processed_events(self, events): - if not isinstance(events, (list, tuple)): - events = [events] - - rows = [] - for event in events: - rows.append(event) - - return self.write_rows(rows) - - def write_rows(self, rows): - if not isinstance(rows, (list, tuple)): - rows = [rows] - - self.dataset.get_writer().write(rows) + return self.write_processed_records(out) diff --git a/tests/snapshots/test_postgres_snapshot.py b/tests/snapshots/test_postgres_snapshot.py index d8ec9e544bd..937a146be52 100644 --- a/tests/snapshots/test_postgres_snapshot.py +++ b/tests/snapshots/test_postgres_snapshot.py @@ -1,15 +1,9 @@ import os # NOQA +import pytest from snuba.snapshots.postgres_snapshot import PostgresSnapshot - -class TestPostgresSnapshot: - - def test_parse_snapshot(self, tmp_path): - snapshot_base = tmp_path / "cdc-snapshot" - snapshot_base.mkdir() - meta = snapshot_base / "metadata.json" - meta.write_text(""" +META_FILE = """ { "snapshot_id": "50a86ad6-b4b7-11e9-a46f-acde48001122", "product": "snuba", @@ -35,21 +29,34 @@ def test_parse_snapshot(self, tmp_path): ], "start_timestamp": 1564703503.682226 } - """) +""" + + +class TestPostgresSnapshot: + + def __prepare_directory(self, tmp_path, table_content): + snapshot_base = tmp_path / "cdc-snapshot" + snapshot_base.mkdir() + meta = snapshot_base / "metadata.json" + meta.write_text(META_FILE) tables_dir = tmp_path / "cdc-snapshot" / "tables" tables_dir.mkdir() groupedmessage = tables_dir / "sentry_groupedmessage.csv" - groupedmessage.write_text( - """id, status -0, 1 -""" - ) + groupedmessage.write_text(table_content) groupassignee = tables_dir / "sentry_groupasignee" groupassignee.write_text( - """id, project_id + """id,project_id """ ) + return snapshot_base + def test_parse_snapshot(self, tmp_path): + snapshot_base = self.__prepare_directory( + tmp_path, + """id,status +0,1 +""" + ) snapshot = PostgresSnapshot.load("snuba", snapshot_base) descriptor = snapshot.get_descriptor() assert descriptor.id == "50a86ad6-b4b7-11e9-a46f-acde48001122" @@ -65,6 +72,20 @@ def test_parse_snapshot(self, tmp_path): assert tables["sentry_groupasignee"] == ["id", "project_id"] with snapshot.get_table_file("sentry_groupedmessage") as table: - content = table.readlines() - assert content[0] == b"id, status\n" - assert content[1] == b"0, 1\n" + line = next(table) + assert line == { + "id": "0", + "status": "1", + } + + def test_parse_invalid_snapshot(self, tmp_path): + snapshot_base = self.__prepare_directory( + tmp_path, + """id +0 +""" + ) + with pytest.raises(ValueError, match=".+sentry_groupedmessage.+status.+"): + snapshot = PostgresSnapshot.load("snuba", snapshot_base) + with snapshot.get_table_file("sentry_groupedmessage") as table: + next(table) diff --git a/tests/test_api.py b/tests/test_api.py index 8e15533b1e3..11b7711d5d6 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -103,7 +103,7 @@ def generate_fizzbuzz_events(self): } } })) - self.write_processed_events(events) + self.write_processed_records(events) def redis_db_size(self): # dbsize could be an integer for a single node cluster or a dictionary @@ -737,7 +737,7 @@ def test_doesnt_select_deletions(self): } result1 = json.loads(self.app.post('/query', data=json.dumps(query)).data) - self.write_processed_events([{ + self.write_processed_records([{ 'event_id': '9' * 32, 'project_id': 1, 'group_id': 1, diff --git a/tests/test_cleanup.py b/tests/test_cleanup.py index e552190f240..303df5611a9 100644 --- a/tests/test_cleanup.py +++ b/tests/test_cleanup.py @@ -16,7 +16,7 @@ def to_monday(d): assert parts == [] # base, 90 retention - self.write_processed_events(self.create_event_for_date(base)) + self.write_processed_records(self.create_event_for_date(base)) parts = cleanup.get_active_partitions(self.clickhouse, self.database, self.table) assert parts == [(to_monday(base), 90)] stale = cleanup.filter_stale_partitions(parts, as_of=base) @@ -24,7 +24,7 @@ def to_monday(d): # -40 days, 90 retention three_weeks_ago = base - timedelta(days=7 * 3) - self.write_processed_events(self.create_event_for_date(three_weeks_ago)) + self.write_processed_records(self.create_event_for_date(three_weeks_ago)) parts = cleanup.get_active_partitions(self.clickhouse, self.database, self.table) assert parts == [(to_monday(three_weeks_ago), 90), (to_monday(base), 90)] stale = cleanup.filter_stale_partitions(parts, as_of=base) @@ -32,7 +32,7 @@ def to_monday(d): # -100 days, 90 retention thirteen_weeks_ago = base - timedelta(days=7 * 13) - self.write_processed_events(self.create_event_for_date(thirteen_weeks_ago)) + self.write_processed_records(self.create_event_for_date(thirteen_weeks_ago)) parts = cleanup.get_active_partitions(self.clickhouse, self.database, self.table) assert parts == [ (to_monday(thirteen_weeks_ago), 90), @@ -44,7 +44,7 @@ def to_monday(d): # -1 week, 30 retention one_week_ago = base - timedelta(days=7) - self.write_processed_events(self.create_event_for_date(one_week_ago, 30)) + self.write_processed_records(self.create_event_for_date(one_week_ago, 30)) parts = cleanup.get_active_partitions(self.clickhouse, self.database, self.table) assert parts == [ (to_monday(thirteen_weeks_ago), 90), @@ -57,7 +57,7 @@ def to_monday(d): # -5 weeks, 30 retention five_weeks_ago = base - timedelta(days=7 * 5) - self.write_processed_events(self.create_event_for_date(five_weeks_ago, 30)) + self.write_processed_records(self.create_event_for_date(five_weeks_ago, 30)) parts = cleanup.get_active_partitions(self.clickhouse, self.database, self.table) assert parts == [ (to_monday(thirteen_weeks_ago), 90), diff --git a/tests/test_groupedmessage.py b/tests/test_groupedmessage.py index fa5a7e79e96..08088dfc606 100644 --- a/tests/test_groupedmessage.py +++ b/tests/test_groupedmessage.py @@ -2,12 +2,19 @@ import simplejson as json from datetime import datetime -from base import BaseTest +from base import BaseDatasetTest +from snuba.clickhouse import ClickhousePool from snuba.consumer import KafkaMessageMetadata -from snuba.datasets.cdc.groupedmessage_processor import GroupedMessageProcessor +from snuba.datasets.cdc.groupedmessage_processor import GroupedMessageProcessor, GroupedMessageRow -class TestGroupedMessageProcessor(BaseTest): +class TestGroupedMessage(BaseDatasetTest): + + def setup_method(self, test_method): + super(TestGroupedMessage, self).setup_method( + test_method, + 'groupedmessage', + ) BEGIN_MSG = '{"event":"begin","xid":2380836}' COMMIT_MSG = '{"event":"commit"}' @@ -94,6 +101,19 @@ def test_messages(self): insert_msg = json.loads(self.INSERT_MSG) ret = processor.process_message(insert_msg, metadata) assert ret[1] == self.PROCESSED + self.write_processed_records(ret[1]) + cp = ClickhousePool() + ret = cp.execute("SELECT * FROM test_groupedmessage_local;") + assert ret[0] == ( + 42, # offset + 0, # deleted + 74, # id + 0, # status + datetime(2019, 6, 19, 6, 46, 28), + datetime(2019, 6, 19, 6, 45, 32), + datetime(2019, 6, 19, 6, 45, 32), + None, + ) update_msg = json.loads(self.UPDATE_MSG) ret = processor.process_message(update_msg, metadata) @@ -102,3 +122,26 @@ def test_messages(self): delete_msg = json.loads(self.DELETE_MSG) ret = processor.process_message(delete_msg, metadata) assert ret[1] == self.DELETED + + def test_bulk_load(self): + row = GroupedMessageRow.from_bulk({ + 'id': '10', + 'status': '0', + 'last_seen': '2019-06-28 17:57:32+00', + 'first_seen': '2019-06-28 06:40:17+00', + 'active_at': '2019-06-28 06:40:17+00', + 'first_release_id': '26', + }) + self.write_processed_records(row.to_clickhouse()) + cp = ClickhousePool() + ret = cp.execute("SELECT * FROM test_groupedmessage_local;") + assert ret[0] == ( + 0, # offset + 0, # deleted + 10, # id + 0, # status + datetime(2019, 6, 28, 17, 57, 32), + datetime(2019, 6, 28, 6, 40, 17), + datetime(2019, 6, 28, 6, 40, 17), + 26, + ) diff --git a/tests/test_optimize.py b/tests/test_optimize.py index 59413a698a6..b91597f2c89 100644 --- a/tests/test_optimize.py +++ b/tests/test_optimize.py @@ -15,25 +15,25 @@ def test(self): base_monday = base - timedelta(days=base.weekday()) # 1 event, 0 unoptimized parts - self.write_processed_events(self.create_event_for_date(base)) + self.write_processed_records(self.create_event_for_date(base)) parts = optimize.get_partitions_to_optimize(self.clickhouse, self.database, self.table) assert parts == [] # 2 events in the same part, 1 unoptimized part - self.write_processed_events(self.create_event_for_date(base)) + self.write_processed_records(self.create_event_for_date(base)) parts = optimize.get_partitions_to_optimize(self.clickhouse, self.database, self.table) assert parts == [(base_monday, 90)] # 3 events in the same part, 1 unoptimized part - self.write_processed_events(self.create_event_for_date(base)) + self.write_processed_records(self.create_event_for_date(base)) parts = optimize.get_partitions_to_optimize(self.clickhouse, self.database, self.table) assert parts == [(base_monday, 90)] # 3 events in one part, 2 in another, 2 unoptimized parts a_month_earlier = base_monday - timedelta(days=31) a_month_earlier_monday = a_month_earlier - timedelta(days=a_month_earlier.weekday()) - self.write_processed_events(self.create_event_for_date(a_month_earlier_monday)) - self.write_processed_events(self.create_event_for_date(a_month_earlier_monday)) + self.write_processed_records(self.create_event_for_date(a_month_earlier_monday)) + self.write_processed_records(self.create_event_for_date(a_month_earlier_monday)) parts = optimize.get_partitions_to_optimize(self.clickhouse, self.database, self.table) assert parts == [(base_monday, 90), (a_month_earlier_monday, 90)]