Class: CassandraObject::Adapters::CassandraAdapter

Inherits:
AbstractAdapter show all
Includes:
CassandraObject::AdapterExtension
Defined in:
lib/initializers/reconnection.rb,
lib/cassandra_object/adapters/cassandra_adapter.rb

Defined Under Namespace

Classes: QueryBuilder

Instance Attribute Summary

Attributes inherited from AbstractAdapter

#config

Instance Method Summary collapse

Methods inherited from AbstractAdapter

#batch, #batching?, #execute_batchable, #initialize, #statement_with_options

Constructor Details

This class inherits a constructor from CassandraObject::Adapters::AbstractAdapter

Instance Method Details

#cassandra_cluster_optionsObject



50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 50

def cassandra_cluster_options
  cluster_options = config.slice(*[
      :auth_provider,
      :client_cert,
      :compression,
      :compressor,
      :connect_timeout,
      :connections_per_local_node,
      :connections_per_remote_node,
      :consistency,
      :credentials,
      :futures_factory,
      :hosts,
      :load_balancing_policy,
      :logger,
      :page_size,
      :passphrase,
      :password,
      :port,
      :private_key,
      :protocol_version,
      :reconnection_policy,
      :retry_policy,
      :schema_refresh_delay,
      :schema_refresh_timeout,
      :server_cert,
      :ssl,
      :timeout,
      :trace,
      :username,
      :heartbeat_interval,
      :idle_timeout
  ])

  {
      load_balancing_policy: 'Cassandra::LoadBalancing::Policies::%s',
      reconnection_policy: 'Cassandra::Reconnection::Policies::%s',
      retry_policy: 'Cassandra::Retry::Policies::%s'
  }.each do |policy_key, class_template|
    params = cluster_options[policy_key]
    if params
      if params.is_a?(Hash)
        cluster_options[policy_key] = (class_template % [params[:policy].classify]).constantize.new(*params[:params]||[])
      else
        cluster_options[policy_key] = (class_template % [params.classify]).constantize.new
      end
    end
  end

  # Setting defaults
  cluster_options.merge!({
                          heartbeat_interval: cluster_options.keys.include?(:heartbeat_interval) ? cluster_options[:heartbeat_interval] : 30,
                          idle_timeout: cluster_options[:idle_timeout] || 60,
                          max_schema_agreement_wait: 1,
                          consistency: cluster_options[:consistency] || :one,
                          protocol_version: cluster_options[:protocol_version] || 3,
                          page_size: cluster_options[:page_size] || 10000
                         })
  cluster_options
end

#cassandra_versionObject



245
246
247
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 245

def cassandra_version
  @cassandra_version ||= execute('select release_version from system.local').rows.first['release_version'].to_f
end

#connectionObject



111
112
113
114
115
116
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 111

def connection
  @connection ||= begin
    cluster = Cassandra.cluster cassandra_cluster_options
    cluster.connect config[:keyspace]
  end
end

#consistencyObject

/SCHEMA



251
252
253
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 251

def consistency
  defined?(@consistency) ? @consistency : nil
end

#consistency=(val) ⇒ Object



255
256
257
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 255

def consistency=(val)
  @consistency = val
end

#create_table(table_name, params = {}) ⇒ Object

SCHEMA



217
218
219
220
221
222
223
224
225
226
227
228
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 217

def create_table(table_name, params = {})
  stmt = "CREATE TABLE #{table_name} "
  if params.any? && !params[:attributes].present?
    raise 'No attributes for the table'
  elsif !params[:attributes].include? 'PRIMARY KEY'
    raise 'No PRIMARY KEY defined'
  end

  stmt += "(#{params[:attributes]})"
  # WITH COMPACT STORAGE
  schema_execute statement_create_with_options(stmt, params[:options]), config[:keyspace]
end

#delete(scope, ids, attributes = {}) ⇒ Object



187
188
189
190
191
192
193
194
195
196
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 187

def delete(scope, ids, attributes = {})
  ids = [ids] if !ids.is_a?(Array)
  statement = "DELETE FROM #{scope.column_family} WHERE #{scope._key} IN (#{ids.map{|id| '?'}.join(',')})"
  arguments = ids
  unless attributes.blank?
    statement += " AND #{attributes.keys.map{ |k| "#{k} = ?" }.join(' AND ')}"
    arguments += attributes.values
  end
  execute(statement, arguments)
end

#delete_single(obj) ⇒ Object



198
199
200
201
202
203
204
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 198

def delete_single(obj)
  keys = obj.class._keys
  wheres = keys.map{ |k| "#{k} = ?" }.join(' AND ')
  arguments = keys.map{ |k| obj.attributes[k] }
  statement = "DELETE FROM #{obj.class.column_family} WHERE #{wheres}"
  execute(statement, arguments)
end

#drop_table(table_name, confirm = false) ⇒ Object



230
231
232
233
234
235
236
237
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 230

def drop_table(table_name, confirm = false)
  count = (schema_execute "SELECT count(*) FROM #{table_name}", config[:keyspace]).rows.first['count']
  if confirm || count == 0
    schema_execute "DROP TABLE #{table_name}", config[:keyspace]
  else
    raise "The table #{table_name} is not empty! If you want to drop it add the option confirm = true"
  end
end

#execute(statement, arguments = []) ⇒ Object



118
119
120
121
122
123
124
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 118

def execute(statement, arguments = [])
  ActiveSupport::Notifications.instrument('cql.cassandra_object', cql: statement) do
    type_hints = []
    arguments.each { |a| type_hints << CassandraObject::Types::TypeHelper.guess_type(a) } unless arguments.nil?
    connection.execute statement, arguments: arguments, type_hints: type_hints, consistency: consistency, page_size: config[:page_size]
  end
end

#execute_async(queries, arguments = []) ⇒ Object



126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 126

def execute_async(queries, arguments = [])
  retries = 0
  futures = queries.map do |q|
    ActiveSupport::Notifications.instrument('cql.cassandra_object', cql: q) do
      connection.execute_async q, arguments: arguments, consistency: consistency, page_size: config[:page_size]
    end
  end
  futures.map do |future|
    begin
      rows = future.get
      rows
    rescue StandardError => e
      retries += 1
      sleep 0.01
      retry if retries <= 3
      raise e
    end
  end
end

#execute_batch(statements) ⇒ Object



206
207
208
209
210
211
212
213
214
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 206

def execute_batch(statements)
  raise 'No can do' if statements.empty?
  batch = connection.batch do |b|
    statements.each do |statement|
      b.add(statement[:query], arguments: statement[:arguments])
    end
  end
  connection.execute(batch, page_size: config[:page_size])
end

#insert(table, id, attributes, ttl = nil) ⇒ Object



158
159
160
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 158

def insert(table, id, attributes, ttl = nil)
  write(table, id, attributes, ttl)
end

#schema_execute(cql, keyspace) ⇒ Object



239
240
241
242
243
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 239

def schema_execute(cql, keyspace)
  schema_db = Cassandra.cluster cassandra_cluster_options
  connection = schema_db.connect keyspace
  connection.execute cql, consistency: consistency
end

#select(scope) ⇒ Object



146
147
148
149
150
151
152
153
154
155
156
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 146

def select(scope)
  queries = QueryBuilder.new(self, scope).to_query_async
  # todo paginate
  arguments = scope.where_values.select.each_with_index{ |_, i| i.odd? }.reject{ |c| c.blank? }
  cql_rows = execute_async(queries, arguments).map{|item| item.rows.map{|x| x}}.flatten!
  cql_rows.each do |cql_row|
    attributes = cql_row.to_hash
    key = attributes.delete(scope._key)
    yield(key, attributes) unless attributes.empty?
  end
end

#statement_create_with_options(stmt, options = '') ⇒ Object



259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 259

def statement_create_with_options(stmt, options = '')
  if !options.nil?
    statement_with_options stmt, options
  else
    # standard
    if cassandra_version < 3
      "#{stmt} WITH bloom_filter_fp_chance = 0.001
        AND caching = '{\"keys\":\"ALL\", \"rows_per_partition\":\"NONE\"}'
        AND comment = ''
        AND compaction = {'min_sstable_size': '52428800', 'class': 'org.apache.cassandra.db.compaction.SizeTieredCompactionStrategy'}
        AND compression = {'chunk_length_kb': '64', 'sstable_compression': 'org.apache.cassandra.io.compress.LZ4Compressor'}
        AND dclocal_read_repair_chance = 0.0
        AND default_time_to_live = 0
        AND gc_grace_seconds = 864000
        AND max_index_interval = 2048
        AND memtable_flush_period_in_ms = 0
        AND min_index_interval = 128
        AND read_repair_chance = 1.0
        AND speculative_retry = 'NONE';"
    else
      "#{stmt} WITH read_repair_chance = 0.0
        AND dclocal_read_repair_chance = 0.1
        AND gc_grace_seconds = 864000
        AND bloom_filter_fp_chance = 0.01
        AND caching = { 'keys' : 'ALL', 'rows_per_partition' : 'NONE' }
        AND comment = ''
        AND compaction = { 'class' : 'org.apache.cassandra.db.compaction.SizeTieredCompactionStrategy', 'max_threshold' : 32, 'min_threshold' : 4 }
        AND compression = { 'chunk_length_in_kb' : 64, 'class' : 'org.apache.cassandra.io.compress.LZ4Compressor' }
        AND default_time_to_live = 0
        AND speculative_retry = '99PERCENTILE'
        AND min_index_interval = 128
        AND max_index_interval = 2048
        AND crc_check_chance = 1.0;
      "
    end
  end
end

#update(table, id, attributes, ttl = nil) ⇒ Object



162
163
164
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 162

def update(table, id, attributes, ttl = nil)
  write_update(table, id, attributes)
end

#write(table, id, attributes, ttl = nil) ⇒ Object



166
167
168
169
170
171
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 166

def write(table, id, attributes, ttl = nil)
  statement = "INSERT INTO #{table} (#{(attributes.keys).join(',')}) VALUES (#{(['?'] * attributes.size).join(',')})"
  statement += " USING TTL #{ttl.to_s}" if ttl.present?
  arguments = attributes.values
  execute(statement, arguments)
end

#write_update(table, id, attributes) ⇒ Object



173
174
175
176
177
178
179
180
181
182
183
184
185
# File 'lib/cassandra_object/adapters/cassandra_adapter.rb', line 173

def write_update(table, id, attributes)
  queries =[]
  # id here is the name of the key of the model
  id_value = attributes[id]
  if (not_nil_attributes = attributes.reject { |key, value| value.nil? }).any?
    statement = "INSERT INTO #{table} (#{(not_nil_attributes.keys).join(',')}) VALUES (#{(['?'] * not_nil_attributes.size).join(',')})"
    queries << {query: statement, arguments: not_nil_attributes.values}
  end
  if (nil_attributes = attributes.select { |key, value| value.nil? }).any?
    queries << {query: "DELETE #{nil_attributes.keys.join(',')} FROM #{table} WHERE #{id} = ?", arguments: [id_value.to_s]}
  end
  execute_batchable(queries)
end