Constructor. :param dict[str, any] cassandra_params: Cassandra parameters :param bytes name: the suffix to be used for the table name :param int buffer_size: the buffer size
(self, cassandra_params, name, buffer_size)
| 367 | QUERY_INSERT = "INSERT INTO {} (key, value, ts) VALUES (?, ?, ?)" |
| 368 | |
| 369 | def __init__(self, cassandra_params, name, buffer_size): |
| 370 | """ |
| 371 | Constructor. |
| 372 | |
| 373 | :param dict[str, any] cassandra_params: Cassandra parameters |
| 374 | :param bytes name: the suffix to be used for the table name |
| 375 | :param int buffer_size: the buffer size |
| 376 | """ |
| 377 | self._buffer_size = buffer_size |
| 378 | self._session = CassandraSharedSession.get_session(**cassandra_params) |
| 379 | # This timestamp generator allows us to sort different values for the same key |
| 380 | self._ts = c_cluster.MonotonicTimestampGenerator() |
| 381 | # Each table (hashtable or key table is handled by a different storage; to increase |
| 382 | # throughput it is possible to share a single buffer so the chances of a flush |
| 383 | # are increased. |
| 384 | if cassandra_params.get("shared_buffer", False): |
| 385 | self._statements_and_parameters = CassandraSharedSession.get_buffer() |
| 386 | self._select_statements_and_parameters_with_decoders = CassandraSharedSession.get_select_buffer() |
| 387 | else: |
| 388 | self._statements_and_parameters = [] |
| 389 | self._select_statements_and_parameters_with_decoders = [] |
| 390 | |
| 391 | # Buckets tables rely on byte strings as keys and normal strings as values. |
| 392 | # Keys tables have normal strings as keys and byte strings as values. |
| 393 | # Since both data types can be reduced to byte strings without loss of data, we use |
| 394 | # only one Cassandra table for both table types (so we can keep one single storage) and |
| 395 | # we specify different encoders/decoders based on the table type. |
| 396 | if b'bucket' in name: |
| 397 | basename, _, ret = name.split(b'_', 2) |
| 398 | name = basename + b'_bucket_' + binascii.hexlify(ret) |
| 399 | self._key_decoder = lambda x: x |
| 400 | self._key_encoder = lambda x: x |
| 401 | self._val_decoder = lambda x: x.decode('utf-8') |
| 402 | self._val_encoder = lambda x: x.encode('utf-8') |
| 403 | else: |
| 404 | self._key_decoder = lambda x: x.decode('utf-8') |
| 405 | self._key_encoder = lambda x: x.encode('utf-8') |
| 406 | self._val_decoder = lambda x: x |
| 407 | self._val_encoder = lambda x: x |
| 408 | table_name = 'lsh_' + name.decode('ascii') |
| 409 | |
| 410 | # Drop the table if are instructed to do so |
| 411 | if cassandra_params.get("drop_tables", False): |
| 412 | self._session.execute(self.QUERY_DROP_TABLE.format(table_name)) |
| 413 | self._session.execute(self.QUERY_CREATE_TABLE.format(table_name)) |
| 414 | |
| 415 | # Prepare all the statements for this table |
| 416 | self._stmt_insert = self._session.prepare(self.QUERY_INSERT.format(table_name)) |
| 417 | self._stmt_upsert = self._session.prepare(self.QUERY_UPSERT.format(table_name)) |
| 418 | self._stmt_get_keys = self._session.prepare(self.QUERY_GET_KEYS.format(table_name)) |
| 419 | self._stmt_get = self._session.prepare(self.QUERY_SELECT.format(table_name)) |
| 420 | self._stmt_get_one = self._session.prepare(self.QUERY_SELECT_ONE.format(table_name)) |
| 421 | self._stmt_get_count = self._session.prepare(self.QUERY_GET_COUNTS.format(table_name)) |
| 422 | self._stmt_delete_key = self._session.prepare(self.QUERY_DELETE_KEY.format(table_name)) |
| 423 | self._stmt_delete_val = self._session.prepare(self.QUERY_DELETE_VAL.format(table_name)) |
| 424 | |
| 425 | @property |
| 426 | def buffer_size(self): |
nothing calls this directly
no test coverage detected