Repository navigation
Expand file tree
/
Copy pathDB.page
More file actions
1953 lines (1733 loc) · 80.1 KB
/
Copy pathDB.page
File metadata and controls
1953 lines (1733 loc) · 80.1 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
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
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
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
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
TODO merge in notes from CS186
Background
----------
resources
- <http://blog.marcua.net/post/117671929/mit-database-systems-6-830-ta-course-notes>
vocab
- PDBMS = MPP
- data cube: hypercube of multi-dim info; analytical aggregation along dims
- a predicate is _sargable_ if DBMS can use index to speed up execution
- usu. requires invertible, order-preserving (monotonic) function on indexed
attribs
fields
- analytical processing, BI, data warehousing
- transaction processing: worry about scaling, consistency
- stream processing, sensor nets
- data integration, deep web, federated db's
- template extraction, schema recognition
- schema resolution
- entity resolution
- data extraction, wrapping (deep web)
- query routing (federated db's)
- information extraction
companies
- BI
- vectorwise: orig monetdb; from CWI, holland; led by marten kersten
- partnered with ingres
- shared-everything (single machine) multicore
- open-source column store (actually row/col hybrid)
- query operators are run via query execution primitives written in
low-level code that allow compilers to produce extremely efficient
processing instrs
- vectors of 100-1000 values get SIMD'd
- founders from CWI
- 3GB/s decompression, 4-5 cycles/tuple, 3-4X ratio
- operate in decompressed domain
- shared scans
- kickfire/C2: entered with TPC-H win
- column store, execution, compression
- FPGA for data- and task-parallelism
- targeting mysql
- died, acquired by teradata at bargain price
- <http://dbmsmusings.blogspot.com/2010/08/thoughts-on-kickfires-apparent-demise.html>
- sybase iq: data mining
- scale-out; C++ UDFs; embedded data mining/predictive analytics lib (modeling & scoring)
- sas: has own internal db
- teradata
- XSPRADA/Algebraix
- paraccel
- kognitio
- enterprisedb: postgresql company
- major contributor to postgresql
- postgresql plus standard server: extensions, eg connectors, geospatial,
replication, caching
- postgresql plus advanced server: deeper oracle compat; infinite cache
- greenplum: postgresql based; acquired by emc
- aster data: customers: myspace; hadoop interop; COTS; investors: sequoia
- ncluster: row-based MPP DBMS based on PG; supports map-reduce
- calpont's infinidb: mysql + columnar storage engine (MPP, dict/token
compression, multicore queries, ACID, no indexes/mat views, "extends" like
netezza's zone maps)
- tried this out since it seemed like most promising free c-store; suffers
from mysql row size limit (no large tables, which is common in OLAP!)
- varchar 8000 max
- infobright: wants multicore queries, ACID
- also column store, competes with infinidb
- tried this out; not the fastest but pretty good sql support out of box
- no CREATE TABLE AS
- luciddb, monetdb: supposed to be slower/less mature than
infinidb/infobright
- netezza: row-based db hw with compression, clustering; acquired by ibm
- compression: prefix compression, huffman coding
- clustering: multi-dimensional simultaneous partitioning w space-filling hilbert curve
- twinfin: well-designed proprietary hw appliance
- twinfini: parallel
- datallegro/ms
- xtremedata: plugs into cpu socket for cpu bus bandwidth
- groovy sql switch: mem oltp rdbms, appliance & sw
- vertica: column store
- stonebraker
- h-store
- main-memory
- ibm soliddb: uses tries instead of btrees for indexes
- oracle timesten
- scaledb: mysql clustering; uses tries instead of btrees for indexes; better
in-mem perf, more compact
- distributed
- continuent tungsten replicator: DB-neutral master/slave async replication
- elephants
- MS: sql server
- shared scans
- oracle/sun: mysql
- oracle exadata: analytical
- postgresql
- hp
- hp neoview: analytical
- ibm
- emc
- marklogic: XML DB with XQuery as primary QL
people
- dave dewitt: uwisc prof; started msr wisc
misc
- XA transactions: open group specification for common dtxn standard
denormalization: simplex/complex reads vs writes
optimizations TODO
- from [ormperf]: caching, outer-join fetching, transactional write-behind, etc.
[ormperf]: http://www.javaperformancetuning.com/news/interview041.shtml
applications
- _on-line transaction processing (OLTP)_: many small updates
- _on-line analytical processing (OLAP)_ aka _data warehousing_: ad-hoc,
long-running queries
- decision support systems
- business intelligence (BI)
- data mining
OLAP
- multi dimensional expressions (MDX): lang for querying OLAP cubes
- resembles SQL on surface; expresses things that may be clumsy in SQL
- basically like pivot tables over RDBMSs
- closely related: cubes, pivot tables, contingency tables
benchmarks
- TPC-C: OLTP
- TPC-E: new OLTP (TODO: does anyone use this? what's different?)
- TPC-H: OLAP
indexes
- primary vs. secondary
- primary: no duplicates, usu. primary key, may store rec directly in index
(otherwise point to rec using rid)
- secondary: support duplicates, points to rec using rid or primary key
- clustered vs. unclustered
- TODO
joins
- bitmap join indices
- build bitmap index over low cardinality outer table R (at least on join
key)
- map into this while scanning over S
- simple hash join
- build hash table of R in memory, and probe into it with S
- on overflow, start spilling to file R'; S similarly spills to S'
- a new hash function h' is used to determine whether incoming R tuples
should be sent to the hash table or to disk
- puzzle: should you just stop reading and perform joins with what you
have so far (i.e. potentially many keys, but not all tuples with that
key are necessarily in the table)? or should you evict things to make
space for all the tuples belonging to each key in the hash table (i.e.
have fewer keys)?
- the former degenerates into O(n^2) nested loop joins - worse because
you can't start to throw out tuples of S at any point
- so you'll need to start making room for the incoming tuples that h'
says should go in memory - but now you need to choose an eviction
policy (probably unimportant since it's never discussed - or I'm just
misunderstanding something)
- recursively (iteratively?) repeat on the overflow files R' and S'
- grace hash join (assume |R| < |S|)
- phase 1: partition R and S into buckets R_i and S_i (all spilled onto disk)
- hash function dictates partitioning; same for R, S so R_i need only be
joined with S_i
- if ~B buffer pages, can have at most ~B partitions (1 output buffer each)
- phase 2: for each R_i, build hash table, then probe S_i against this
- if R_i is too big and doesn't fit in memory, overflow
- can recursively partition R_i back onto disk into sub-partitions
- hybrid hash join
- uses same partitioning, but instead of using all of memory for output
buffers, use part of it for the first bucket's hashtable (i.e. build hash
table for bucket 1 immediately)
- usually you want to make sure you have enough buckets so that each bucket
will subsequently be able to fit in memory; if this is the primary goal,
then the immediate hash-tabling of bucket 1 is opportunistic (using
whatever pages remain)
- makes sense when you have larger memory (N buffer pages >> B buckets)
logging and recovery
- buffer management policies
- _steal_: steal frames with dirty pages
- thus dirty pages are written ("swapped") to disk
- faster but complicates recovery
- _force_: force all updates to disk before commit
- thus committed pages might not be on disk
- slower but helps recovery
no steal steal
-------------------- -------------------- --------------------
no force no undo, redo undo, redo
force no undo, no redo undo, no redo
- write-ahead logging (WAL)
- undo info: force log record to disk for an (uncommitted) update *before*
corresponding data page is written to disk (swapped/stolen)
- redo info: force log record to disk for a xact before commit
- log seq num (LSN): increasing ID for each log record
- each data page contains a _page LSN_, the LSN of the last log record for an
update to that page
- system keeps track of _flushed LSN_, the max LSN flushed so far
- technique to ensure coherent page written to disk: include in commit marker
the CRC of the data
- add'l barriers vs CRCs: <http://lwn.net/Articles/283161/>
- WAL: before page $i$ written to disk, log must satisfy pageLSN $\le$
flushedLSN
- log records
- fields: LSN, prevLSN, XID, type, [pageID, len, offset, beforeImg, afterImg]
- prevLSN: LSN of last record *in this xact*
- type: update, commit, abort, checkpoint (for log maintenance), compensation
(for undo), end (end of commit/abort)
- [...]: updates only
- in-memory tables
- xact table: XID, status (running/committing/aborting), lastLSN (last
written by xact)
- dirty page table: pageID (entry exists for each dirty page), recLSN (log
record which first dirtied the page)
- checkpoint periodically for faster recovery
B+ trees TODO
- B-link-trees: sibling links are for concurrency, as a "temp repair bridge"
while parent nodes are being rebalanced
LSM trees
- uses _log-structured merge (LSM)_ trees to create full DB replicas using
purely sequential IO
- LSM-trees don't become fragmented; provide fast, predictable index scans
- LSM-tree lookup perf comparable to B-tree lookups
- good for many (amortized) random writes; bad for random reads (need to hit
multiple _components_)
mysql
- engines
- myisam
- innodb
- percona's xtradb: innodb replacement that scales better to mcores and
better utilizes mem
- mariadb's maria: evolve myisam toward innodb; adds txns, crash safety, etc
- forks
- mariadb: founder's community-oriented fork; dissatisfied with sun procedure
- drizzle: complete refactoring/slimming by brian aker, long-time mysql dev
innodb
- index-organized tables
- full row data in PK index; saves space
- 2ary indexes refer to PK value, not phys loc (can change due to splitting)
- can't add/drop/change PK def'n without full-table rewrite
- can't scan in physical page order; must scan in key order
- supports index-only scans; covering indexes run faster
- MVCC
- updates copy old row to rollback segment first
- deletes mark row for deletion, insert info into rollback segment too
- "purge" GC can quickly determine what to discard from rollback segment
- purge is single-threaded
- not sure how deleted space is reclaimed/reused/kept around (doesn't info
about it get removed from the rollback segment by the purge process?)
sqlite
- concurrency
- one global lock
- multi-process concurrency: only whole-database locks
- unlocked, shared, reserved, pending, exclusive
- reserved: intend to write (upgrade to exclusive)
- pending: waiting to upgrade to exclusive; prevents new locks
- isolation levels: eg `BEGIN DEFERRED`
- deferred: acquire only on first query
- immediate: acquire reserved lock right away; reads can still proceed
- exclusive: acquire exclusive lock right away
- shared-cache concurrency (disabled by default)
- supports multiple cnxns to shared cache; 1 txn per cnxn at a time
- 3 levels of locking
- txn-level locking: at most 1 write txn at a time; txn implicitly
becomes writer on first INSERT/UPDATE/DELETE; coexist with any readers
- table-level locking: reader and writer locks
- schema-level locking: on sqlite_master catalog table; some special
rules
- isolation levels
- default isolation level is (actually) serializable, since <=1 writer
per db at a time
- read uncommitted: cnxn doesn't take read locks; no blocking
- features
- heavy testing: 100% branch coverage; fault injection (malloc, IO);
virtual file system snapshots to simulate crashes
- virtual table: you can export a table interface and provide callback
implementations for events on this table
- storage
- data stored directly in btree
sql server 2008
- integration services: ETL tools
- reporting services
- analysis services: BI tools
- powerpivot for excel/sharepoint makes report designing easier/interactive
- full text search engine
- r2
- streaminsight: CEP
- master data management: dealing with multiple DBs/sources
- reporting services integration w sharepoint CMS: publish olap analytics
postgresql
- table data and PK index are separate
- can scan in physical page order; fast
- supports changing PKs wo full-table rewrite; even allows concurrent writes
- no index-only scans yet; covering indexes don't help
- adding more cols to indexes slows things down
- MVCC: on update, old row left in place, new row quickly inserted
- VACUUM must scan all data as a result
- VACUUM is multi-threaded
- users: disqus, skype, debian, omniti
- <http://www.postgresql.org/about/users>
postgresql vs mysql
- table org: PK index & table data sep vs same (see sections)
- neither can scan index in physical order, only in key order
- <http://rhaas.blogspot.com/2010/11/mysql-vs-postgresql-part-1-table.html>
- MVCC GC: VACUUM vs purge
- innodb: updates must copy old row to rollback segment
- PG: faster updates, slower VACUUM scans all data
- <http://rhaas.blogspot.com/2011/02/mysql-vs-postgresql-part-2-vacuum-vs.html>
oracle
- table data and PK index can be same (mysql) or sep (pg)
distributed postgresql
- postgres-r: async master-master replication
- pgcluster: sync multi-master replication
- pgpool-II: sync
distributed mysql
- mysql cluster/NDB (network database)
- node types
- mgmt node: config server, start/stop nodes, run backup, etc
- data node: number of data nodes = replicas * fragments
- sql node or api node
- database
- partition: a fragment of the database; these are replicated
- replica: a copy of a partition
- hash-partitioned on table primary keys
- stored in mem or on disk asynchronously
- organization
- data node: manages a replica
- node group: manages a partition
- checkpoints
- local: save all data to disk; every few mins
- global: txns for all nodes synced; redo log flushed; every few secs
- system
- replication may be sync or async
- 2PC
- no referential integrity
- no txn isolation above read committed, so dtxns are naive
- replication via mysql replication (replicate entire cluster)
- mysql replication
- system
- async replication or semi-sync replication (wait for receipt; contributed
by google in 2007)
- supports multi-master
- simple "greatest timestamp wins" conflict resolution
- no auto failover
- formats
- statement-based replication: stream the queries
- row-based replication: stream the updates
- mixed format: default; usu statements, but rows for things like UUIDs
- mysql multi-master replication manager (MMM)
- async replication *with* auto failover
- also: flipper, tungsten
clustering
----------
blenddb
- colocate related tuples on same pages in navigational queries
- eg web page for a movie: `actors >< movies >< directors; movies -- ratings`
- eg further navigate into
- construct _paths_ starting from prim key into one table and following
foreign key joins to other tables; eg `movie.id -> acted.movieid ->
actor.id`, `movie.id -> directed.movieid -> director.id`, `movie.id ->
ratings.movieid`
- pull in all tuples read by those paths
- clustering is an offline process
- optimization i didn't try to fully understand: prioritize path heads by
how frequently they spill over one page (don't waste space on path sets
that are too big anyway and require multiple seeks?)
- faster than performing joins online and less space than materialized views
- related work
- (stochastic) clustering in OODBMSs: could use their algos here, but they're
about migration not replication
- _merged indices_ of multiple cols/tables: can cluster on this, but that
only benefits joins on a single key, i.e. no navigation
- memcached: blenddb is not a cache; prepped for any data
- multi-table indices, eg oracle clusters: crappy implementation of merged
indices that requires sep keys to be on sep pages
- first work on clustering in RDBMS
RDBMSs
------
amazon relational database as a service (RDS)
- mysql 5.1
- synchronous wan repl (multi-az)
- TODO how does it scale?
OLTP Through the Looking Glass, and What We Found There (stavros, dna, madden, stonenbraker)
- lots of CPU overhead in conventional RDBMS
- removed features from Shore to get to main memory DB
- Shore uses user level non-preemptive threads with IO processes
- [hence latching overheads not from multiprocessing; unclear why latching is
so impactful or even prevalent]
- went from 640 to 12700 TPS; down to 6.8% of orig work
- removed in order:
- 35% buffer mgr
- 14% latching: many structs
- 16% locking: 2PL lock mgr
- 12% WAL: building records; pinning records; IO
- 16% hand-coded optimizations: tuned B trees
- "every reason to believe that many OLTP systems already fit or will soon fit
into main memory"
- [but no cost benefit analysis; james hamilton says even move to SSDs is not
worth it, contrary to myspace's claims]
scientific databases
--------------------
general
- projects
- extremely large database (xldb): workshop started by jacek becla of slac
- 55 PB raw images, 100 PB astronomical data for LSST
- involvement
- ppl: stonebraker, dewitt, kersten
- vendors: teradata, greenplum, cloudera
- companies: ebay, web companies
- scidb: stonebraker, dewitt
- A data model based on multidimensional arrays, not sets of tuples
- A storage model based on versions and not update in place
- Built-in support for provenance (lineage), workflows, and uncertainty
- Scalability to 100s of petabytes and 1,000s of nodes with high degrees of
tolerance to failures (including query fault tolerance)
- Support for "external" objects so that data sets can be queried and
manipulated without ever having to be loaded into the database
- Open source in order to foster a community of contributors and to insure
that data is never "locked up"---a critical requirement for scientists
- refs
- <http://www.dbms2.com/2009/10/03/issues-in-scientific-data-management/>
- <http://www.dbms2.com/2009/09/12/xldb-scid/>
scientifica (DB group meeting talk, Mike Stonebraker, 4/24/08)
- multi-dimensional array-based
- interesting question is storage manager: how to partition the data (and
then in each partition, how to chunk it)
- uncertainty
- minimal, since error analysis is different for each application
- use error bars
- lineage, provenance
- need also to remember the derivation process, not just origin
- interest from various scientific institutions but also amazon, google, etc.
- CERN had planned to use Objectivity
Efficient Provenance Storage
- main problem: compressing binary xml
- intro
- fields: science, law, etc.
- dataset 270mb, provenance store 6gb
parallel/distributed olap
-------------------------
"query evaluation techniques for large databases" TODO
HARBOR (Lau, VLDB06)
- intro
- integrate mechanisms for recovery and high availability
- take advantage of data redundancy in replicated dbms
- features
- no logging
- recover without quiescing
- traditional approaches
- logging
- approach overview
- data warehousing techniques (assumptions)
- data replication: logical (phys unnec); helps perf anyway
- snapshot isolation
- large ad-hoc query workloads over large read data sets intermingled
with a smaller number of OLTP transactions
- using snapshot isolation avoids lock contention (TODO is this a
requirement?)
- harbor uses time travel mechanism similar to snapshot isolation
- historical queries of past can proceed without locks since past data
never changes
- fault tolerance model
- k-safety: requires k+1, k may fail
- _no network partitions_, no BFT
- historical queries
- assume T < now, so no locks needed (result set don't change due to
updates/deletes after T)
- recovery approach
- checkpoints: flush dirty pages and record current time T
- after coming back up, execute historical query on other live sites that
replicate your data; this catches you up to a point closer to present
- execute standard non-hist query to catch up to present; this uses read
lock, but is shorter than the prev query
- query execution
- recovery query is a range query on T, so break up data into time-partition
segments; means normal queries will need to merge from all the segments
- commit processing
- 2pc with changes
- commits include commit times (for modified tuples)
- in-mem lists of modified tuple ids can be deleted on commit/abort
- opt 2pc: no need for logging/disk writes, except coord
- opt 3pc: no need for logging/disk writes, including coord (persist state
onto workers, basically); non-blocking (can establish timeout)
- recovery
- query remote, online sites for missing updates
- evaluation
- compare against 2pc and aries
A Performance Evaluation of Four Parallel Join Algorithms in a Shared-Nothing
Multiprocessor Environment (DeWitt, SIGMOD '89)
- tried out some hash joins in gamma
- simple hash join
- partition tuples among the join nodes (via a _split table_)
- overflow possible at any of the join nodes; overflow handling is same as in
centralized version, i.e., each node has its own overflow h' and R'
- the split table is augmented with h', so S tuples are directed to either
the correct join node, or to S'
- grace hash join
- each partition (bucket) is in turn partitioned across all _partition
nodes_
- then, for each bucket, its tuples are partitioned across the _join nodes_,
where the hash table is constructed
- hence, all partition nodes talk to all join nodes
- hybrid hash join
- the bucket-1 tuples go straight to some join nodes, and the rest are
spilled to the partition nodes
- it seems in their system that join nodes are diskless, so that's all they
can do - but if the set of join nodes overlaps with the set of partition
nodes, then you can imagine immediately processing a number of buckets
immediately (as many as there are nodes - the immediate hash table will
just consume part of the memory of each of the join nodes)
chained declustering (dewitt 92)
- problem space: distributed data placement (replication, clustering, etc.)
strategies for high availability, reliability in shared nothing dbs
- compared against:
- tandem: mirrored disks
- each partition mirrored on 2 disks that are each connected to same 2
controllers that are each connected to same 2 cpu's
- on cpu failure, cpu 2 has to support the load of parts 1 & 2
- [silly: the problem with this comparison is that it's comparing a
shared-nothing system against a shared-something system]
- equivalent shared nothing system is one where you treat each disk pair
as the single partition they're storing, and replicate each such
partition onto 2 machines
- and that would be obv. stupid
- teradata
- each partition has primary and backup copies
- in $n$-node system, fragment the backup copy into $n-1$ pieces scattered
over the other nodes
- good innate load redistribution on failure
- but expensive writes in normal operation
- RAID5: maintain (parity) bytes
- even a blind write of one sector of a block requires 4 accesses (1 extra
"round"): read all 3 other disks then write parity
- chained declustering
- eg
- node 1: part 1 primary + part $n$ backup
- node 2: part 2 primary + part 1 backup
- node 3: part 3 primary + part 2 backup
- etc
- on failure, evenly redistribute the partition over all other hosts
- in larger systems, don't spread across all nodes; partition up system
(called "relation clusters")
- [this was not designed for high scaling distributed systems]
- assuming indep failures
- back to inefficient system if nodes $i$ and $i+2$ fail
- compared against weird old architectures, not giant commodity pc clusters
in datacenters
hybrid-range partitioning strategy: a new declustering strategy for
multiprocessor database machines (dewitt, vldb90)
- partitionings in gamma: round-robin, range, hash
- for small range queries, range partitioning can localize execution to only
relevant processors and is faster
- for large range queries consuming significant CPU/IO resources, hash/RR can
parallelize, lowering response time
- hybrid-range partitioning strategy (HRPS): do query analysis to get optimal
degree of intra-query parallelism, balancing range and hash/RR
- fragments contain $FC$ tuples and fragments contain unique range of values
of partitioning attr
- collect some stats on avg query characteristics, plugging into formulas for
first derivatives to optimize (e.g.) # servers to distributed across
- The hybrid-range partitioning strategy is an alternative to hash-
and range-partitioning that uses query analysis to compute the typical
resource consumption requirements of queries and from this determines the
optimal partition size for range-partitioning. The scheme attempts to reduce
resource consumption by localizing small range queries while declustering and
parallelizing long-running range queries; however, HRPS is only applicable to
range queries.
adaptively parallelizing distributed range queries (adam silberstein, brian
cooper, vldb09)
- problem space: parallelizing range queries
- traditionally, lay out data to achieve highest throughput/max parallelism
- in high-scaling systems like PNUTS/bigtable, maximizing server parallelism
produces too much throughput for 1 client
- assumption: tables are range-partitioned
- proposal: adaptive approach
- step 1: adaptive server allocation
- find ideal parallelism for single query execution
- depends on: query selectivity, client load, client-server BW, etc
- min # servers $K_q$ for query $q$ is adjusted as scanning proceeds
- step 2: multi-query scheduling
- even after knowing $K_q$, must avoid contention (assigning too many
queries to same server simultaneously)
- num concurrent scans on a server to avoid random IO can result in faster
execution for all queries
- minimize disk contention and ensure all queries get good perf
- max # queries $L_s$ on server $s$
- tunable: can favor short vs long queries or high vs low priority
- implemented & evaluated in PNUTS
gamma (project dewitt led 8x-92)
- shared-nothing dbms; horizontal partitioning
- implemented and evaluated parallel joins: simple hash, grace hash, hybrid
hash, sort-merge
sharing
- shared memory
- shared disk
- shared nothing
- the first two are relics anyway
why main memory is bad for olap (DB group meeting talk, Mike Stonebraker,
4/24/08)
- compared: Vertica with disk, Vertica with ramdisk, and Vertica WOS only
- seek times disappear; a few disks can satisfy bandwidth (even for several
cores)
- absence of indexes contributes 2x
- compression contributes 7x
- SAP will compare their TREX against Vertica
distributed dbms
----------------
building a database on s3
- this system completely gives up concurrency control (write-write conflicts)
and coherency (write-read consistency)
- durability via SQS
- highly inefficient, but claims "good for internet speeds"; several seconds
per transaction
- <http://www.dbis.ethz.ch/people/kraskat/SIGMOD-s3-slides.pdf>
rose: compressed, log-structured replication (russell sears, VLDB08)
- ROSE: replication oriented storage engine
- storage engine for high-tput replication
- targets seek-limited, write-intensive workloads that do near-real-time OLAP
- target usage in replication: replication log (of actual data) comes
streaming into replicas, which are fed as input to ROSE (which itself just
uses STASIS for logging LSM tree ops)
- although targeting replication, provides high-tput writes to any app that
updates tuples without reading existing data, eg append-only, streaming,
and versioning DBMSs
- write perf depends on replicas' ability to do writes without looking up old
values; reads are expensive (esp. of deleted values)
- maintains multiple versions/snapshots for each item, supporting eg OCC, MVCC,
and more
- deletion: tombstone insertion
- introduce page compression format that takes advantage of LSM-tree's
sequential, sorted data layout
- increases replication tput by reducing sequential IO
- enables efficient tree lookups by supporting small page sizes & doubling as
an index of the values it stores
- any scheme that can compress data in a single pass and provide random
access to compressed values could be used by rose
- replication envs have multiple readers, 1 writer; thus rose gets atomicity,
consistency, isolation to concurrent txns without resorting to rollbacks,
blocking index reqs or interfering with maintenance tasks
- rose avoids random IO during replication and scans
- leaves more IO capacity for queries than existing systems
- provides scalable, real-time replication of seek-bound workloads
- analytical models and experiments: OOMs greater replication bandwidth
- column-store within pages (a la PAX)
- within each column, doesn't sort by column values; just stores in tuple
slot id order
- adds overhead to lookups
- appends: append to each col; buffered staging space
OLTP
----
the end of an architectural era (it's time for a complete rewrite) (stonebraker, madden, dna, stavros, nabil, pat helland)
- features
- main memory
- "vast majority of OLTPs <1 TB in size, and growing slowly"
- [implies you're throughput bound vs. capacity bound]
- timesten, soliddb: inherit RDBMS baggage of System R, incl. disk-based
recovery log, dynamic locking
- single-threading vs. multi-threading and resource control (sync)
- OLTP txns lightweight
- current systems interleave CPU and IO via synchronization
- elasticity vs. fork-lift upgrades
- high avail through shared-nothing replication
- no knobs
- txn and schema characteristics
- OLTP tend to have highly partitionable tree schemas [think entity groups]
- [unsubstantiated; I do believe a more cyclic reasoning: high-scaling apps
already are partitionable]
- _constrained-tree application_: txns run within 1 partition
- _one-shot_ txn: can be executed in parallel without requiring intermed
results to be serialized; still needs 2PC; can be done via vert part
- also: _two-phase_, _strongly two-phase_, commutativity
- system architecture
- all txns are stored procedures; runtime in same thread as DB
- distributed txn mgmt: single-sited, one-shot, general
- replication and recovery
- DB designer: make as many txn classes as possible single sited
- [this is focused on static property of txn classes; I would like to focus
on dynamic property of txn instances, i.e. individual objects]
- [lot of faith in TPC-C]
h-store concurrency control (evan jones)
- meat is speculative execution with multi-partition tweak
- local speculation scenario
- say you have the following txns, where x=2 is on partition p and y=10 is on q:
- txn A: swap x, y
- this is a multi-partition txn
- this is multi-fragment on p
- txn B1: x++
- txn B2: x++
- end result should be x=12, y=2 OR x=11, y=3 OR x=10, y=4
- timeline: say we execute in A, B1, B2 order
- read x and read y, sending to coord
- cannot yet begin speculation of B1/B2, since A left inconsistent state
- ow, second fragment of A will overwrite B1/B2's effects
- end result will be x=10, y=2; not serializable
- locking or occ can "execute" B1, but they have overhead, and in this
(conflicting) case, they will be blocked
- each partition gets second fragment of A from coord
- each partition sends ack for A, waits to commit (2PC)
- while waiting, can start speculating B1/B2
- must queue results; release only once no longer speculative
- first fragment of a multi-part txn can always be speculated like this
- can do better for speculating these first fragments of multi-part txns
if using central coord and coord is aware of speculative txns; see next
sec
- must keep undo buffers; must undo/re-exec B1/B2 *iff* A aborts
- if you have one coord for a batch of multi-partition txns...
- A, B1, C, B2, where C increments both x,y; note C is multi-part single-frag
- everything goes as before up till B1 speculation
- because A and C have same coordinator, q can return the result to its
portion of C, along with indication that result is speculating on A
- this is the *end* of all of C; it's a speculative ack
- after committing A, coord can immediately commit C
- p can speculate its fragment of C, but like B1, cannot send this result,
because coord is unaware of prior B1 txn, since single-partition txns do't
go through the coord
- where does this help?
- allows a sequence of multi-part txns, each composed of a single frag at
each part, to be executed without blocking; call these _simple multi-part
txns_
- these are common, e.g. updates to highly replicated read-mostly tables
- e.g. updates based on a non-partition column need to search all
partitions
- all distributed txns in TPC-C are simple multi-part txns
- limitations
- only applicable after executing last fragment of a multi-part txns, so
not useful for non-simple txns
- require same coord; on abort, need to be able to cascade aborts
- might be able to avert this by assigning txns a global order
- speculation is on the success of a 2PC
- unlike concurrency control, it assumes all txns conflict
voltdb
- 1.0
- server-compiled java stored procedures
- specify partitioning columns for each table
- $k$-safety; not viewstamped replication; static config
- supports DB-wide replicated (read-mostly) tables
- no concurrency control; strictly sequential operation
- statements can't do joins on 2 partitioned tables
- supports aggs
- periodic/continuous snapshots to disk
- 2PC for dtxns?
- early beta w 150 customers featured bunch of web gaming companies and bunch
of use cases resembling stream processing
percolator (daniel peng, frank dabek, google, osdi10)
- btwn DBMS & mapreduce: online processing that tolerates latency
- latency comes from conflicts
- built on bigtable, chubby
- no blocking - on conflict, abort
- locks taken at the sites
- one of the locked sites designated as primary
- ACID snapshot isolation (write-write conflicts only; reads are not monitored)
- centralized timestamp oracle hands out timestamps; can reply to 2B req/s
- allows others to run their own systems
- solves blockage on 2PC coordinator failure: after some time, try to undo the
primary lock
- the primary lock site is "the" commit site
- bigtable has row txns, which allows for atomic test-and-set to rollback or
commit txns at the primary lock site
- abort & backoff if see a lock w any timestamp or a write w more recent
timestamp than current txn timestamp
- _observers_ run on storage servers; multiple changes may coalesce into one
observer run
- implemented via _notifications_ in special bigtable column group/column
- slower than MR due to excessive RPCs, but major engineering gains
- lock operations made faster by adding special RPC to bigtable
- app: incremental indexing
haystack (facebook, osdi10)
- facebook photo store
- replaces NFS NAS solution that choked on sheer # files and excessive disk
seeks
- photos stored in 100GB append-only physical volumes (files) that are replicas
of logical volumes
- lookup operation
1. get metadata/location-specifying URL from the metadata server (haystack
directory), which balances load across physical volumes
2. browser issues img requests to the CDN, which fwds to haystack, or
directly to haystack
3. lookup in haystack cache (DHT, basically a CDN)
4. haystack store looks up some position info in in-memory index and grabs
the photo in 1 seek
5. put into haystack cache iff requesting directly (CDN can cache) &
requesting from a write-enabled store (recently-written photos are most-read)
- index has filename, offset, and size for each photo ID
- photo IDs have random number (_cookie_) to prevent guessing valid URLs
- writes synchronously append to volume and asynchronously write index
checkpoints for faster recovery (the volume is like the DB log)
- batched uploads when possible; luckily many users upload whole albums
- updates: latest ver in physical offset has highest offset; if diff physical
vol, must update
- deletes synchronously update physical volumes & index; compaction does GC
- compactions: processes deletes & duplicate keys
- occasionally, full reset needed of a node, which takes many hours
- use XFS extent-based FS bc it has efficient prealloc & blockmaps that fit in
mem
Data Integration
----------------
Indexing Dataspaces (Xin Dong, Alon Halevy, sigmod07)
- INTRO
- users explore and use predicate queries and neighborhood
- keyword queries
- traditionally: inverted list or index tree or XML
- PROBLEM DEFINITION
- INDEXING HETEROGENEOUS DATA
- eg: bibtex, word, ppt, email, www, xls, xml, rdb
- triples: (instance, attr, val) or (inst, assoc, inst)
(directional)
- how to extract from data sources to this model is a separate
topic
- synonyms among attrs & assocs
- hierarchies: 2 types (don't distinguish between them)
- sub-property: eg father < parent, contactAuthor < author
- sub-field: eg city < addr, firstName < name
- [not about complex structures' components, but see 4.4]
- [what happens when name is a unique inst? question of
extraction]
- QUERYING HETEROGENEOUS DATA
- predicate queries
- attr pred: matches on direct attr/sub-attr
- assoc preds: matches instances that have a (sub-)association
with another instance that has any matching attr
- [sub-attr/sub-assoc: either sub-property or sub-attr]
- [(firstName 'Yang') won't match an obj with name 'Yang']
- neighborhood keyword queries
- relevant inst: has direct attr
- associated inst: assoc with a relevant inst
- user need not have knowledge of exact schema
- ranking is a separate issue
- [both queries limited to 2 degrees]
- INVERTED LISTS
- store occurrence counts
- INDEXING STRUCTURE
- INDEXING ATTRS
- attr inverted lists (ATIL): inverted list on (keyword//attr//)
- INDEXING ASSOCS
- attr-assoc inverted list (AAIL): ATIL + (keyword//assoc//)
entries too
- [doesn't count occurrences within each assoc inst]
- easy extension to support k-ary assocs
- INDEXING HIERARCHIES
- INDEX WITH DUPLICATION
- dup-ATIL: duplicate (eg Yang//firstName//, Yang//name//)
- increases index size
- INDEX WITH HIERARCHY PATH
- hier-ATIL: eg Yang//name//firstName//
- use row-wise deltas to minimize space overhead
- prefix serach more expensive
- HYBRID INDEX
- summary rows shadow other rows and end in additional //
- [need to skip the shadowed rows somehow - easy problem]
- start with hier-ATIL; add summary row for prefix p if count of
other rows with prefix p exceeds threshold
- [what if: a/a/{a,b,c/{a,b,c}}?]
- SCHEMA-LEVEL SYNONYMS
- assoc in one src = attr in another: naturally handled
- term heterogeneity: have table mapping sets of synonyms to
canonical names, and use this replacement in index and queries
- NEIGHBORHOOD KEYWORD QUERIES
- keyword inverted list (KIL): hybrid-AAIL, but also summarize
prefixes that correspond directly to keywords, i.e. k//
- EXPERIMENTAL EVALUATION
functiondb - arvindt
bootstrapping pay as you go - alon halevy
pay as you go - alon halevy
oltp ttlg
uncovering the relational web
column-stores vs. row-stores
column stores
-------------
general
- column stores only recently gained traction because there's enough memory for
fetching large ranges of blocks
hybrid approaches
- PAX (Weaving Relations for Cache Performance by Natassa Ailamaki, David
DeWitt, Mark Hill, and Marios Skounakis, VLDB 2001)
- store data by column *within disk block*
- pros
- CPU efficiency of C-stores
- improved cache hit ratio & mem bandwidth
- better compression
- easy to implement in row-store
- cons
- IO properties of row-stores
- cache prefetching on modern CPUs makes PAX cache efficiency obsolete
- fractured mirrors (A Case for Fractured Mirrors by Ravi Ramamurthy, David DeWitt, and Qi Su, VLDB 2002)
- replicate multiple ways; c-stores are OOM faster than row-stores and
vice-versa
- cons
- may need higher degrees of replication
- harder impl: need full c-store & row-store exec engines & query
optimizers
- writes can be a problem
- fine-grained hybrids (Cudre-Mauroux, Wu, and Madden, CIDR 2009)
- individual tables can be row- or col-oriented; declarative storage
- pros: get perf advantages of fractured mirrors without additional