MCPcopy Create free account
hub / github.com/BIT-DataLab/LakeBench / __init__

Method __init__

join/LSH/datasketch/storage.py:369–423  ·  view source on GitHub ↗

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)

Source from the content-addressed store, hash-verified

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):

Callers

nothing calls this directly

Calls 4

get_sessionMethod · 0.80
get_bufferMethod · 0.80
get_select_bufferMethod · 0.80
getMethod · 0.45

Tested by

no test coverage detected