Principles of Distributed Database
Systems
M. Tamer Özsu
Patrick Valduriez
© 2020, M.T. Özsu & P. Valduriez 1
Outline
Introduction
Distributed and Parallel Database Design
Distributed Data Control
Distributed Query Processing
Distributed Transaction Processing
Data Replication
Database Integration – Multidatabase Systems
Parallel Database Systems
Peer-to-Peer Data Management
Big Data Processing
NoSQL, NewSQL and Polystores
Web Data Management
© 2020, M.T. Özsu & P. Valduriez 2
Outline
Distributed Data Control
View management
Data security
Semantic integrity control
© 2020, M.T. Özsu & P. Valduriez 3
Semantic Data Control
Involves:
View management
Security control
Integrity control
Objective :
Ensure that authorized users perform correct operations on the
database, contributing to the maintenance of the database
integrity.
© 2020, M.T. Özsu & P. Valduriez 4
Outline
Distributed Data Control
View management
Data security
Semantic integrity control
© 2020, M.T. Özsu & P. Valduriez 5
View Management
View – virtual relation EMP
generated from base relation(s) by a ENO ENAME TITLE
query
E1 J. Doe Elect. Eng
not stored as base relations E2 M. Smith Syst. Anal.
E3 A. Lee Mech. Eng.
Example: E4 J. Miller Programmer
CREATE VIEW SYSAN(ENO,ENAME) E5 B. Casey Syst. Anal.
E6 L. Chu Elect. Eng.
AS SELECT ENO,ENAME E7 R. Davis Mech. Eng.
E8 J. Jones Syst. Anal.
FROM EMP
WHERE TITLE= "Syst.
Anal."
© 2020, M.T. Özsu & P. Valduriez 6
View Management
Views can be manipulated as base relations
Example :
SELECT ENAME, PNO, RESP
FROM SYSAN, ASG
WHERE [Link] = [Link]
© 2020, M.T. Özsu & P. Valduriez 7
Query Modification
Queries expressed on views
Queries expressed on base relations
Example :
SELECT ENAME, PNO, RESP
FROM SYSAN, ASG
WHERE [Link] = [Link]
SELECT ENAME,PNO,RESP
FROM EMP, ASG
WHERE [Link] = [Link]
AND TITLE = "Syst. Anal."
© 2020, M.T. Özsu & P. Valduriez 8
View Management
To restrict access
CREATE VIEW ESAME
AS SELECT *
FROM EMP E1, EMP E2
WHERE [Link] = [Link]
AND [Link] = USER
Query
SELECT *
FROM ESAME
© 2020, M.T. Özsu & P. Valduriez 9
View Updates
Updatable
CREATE VIEW SYSAN(ENO,ENAME)
AS SELECT ENO,ENAME
FROM EMP
WHERE TITLE="Syst. Anal."
Non-updatable
CREATE VIEW EG(ENAME,RESP)
AS SELECT ENAME,RESP
FROM EMP, ASG
WHERE [Link]=[Link]
© 2020, M.T. Özsu & P. Valduriez 10
View Management in Distributed DBMS
Views might be derived from fragments.
View definition storage should be treated as database
storage
Query modification results in a distributed query
View evaluations might be costly if base relations are
distributed
Use materialized views
A materialized view is a duplicate data table created by combining data
from multiple existing tables for faster data retrieval.
© 2020, M.T. Özsu & P. Valduriez 11
Materialized View
Origin: snapshot in the 1980’s
Static copy of the view, avoid view derivation for each query
But periodic recomputing of the view may be expensive
Actual version of a view
Stored as a database relation, possibly with indices
Used much in practice
DDBMS: No need to access remote, base relations
Data warehouse: to speed up OLAP
Use aggregate (SUM, COUNT, etc.) and GROUP BY
© 2020, M.T. Özsu & P. Valduriez 12
Materialized View Maintenance
Process of updating (refreshing) the view to reflect
changes to base data
Resembles data replication but there are differences
View expressions typically more complex
Replication configurations more general
View maintenance policy to specify:
When to refresh
How to refresh
© 2020, M.T. Özsu & P. Valduriez 13
When to Refresh a View
Immediate mode
As part of the updating transaction, e.g. through 2PC
View always consistent with base data and fast queries
But increased transaction time to update base data
Deferred mode (preferred in practice)
Through separate refresh transactions
No penalty on the updating transactions
Triggered at different times with different trade-offs
Lazily: just before evaluating a query on the view
Periodically: every hour, every day, etc.
Forcedly: after a number of predefined updates
© 2020, M.T. Özsu & P. Valduriez 14
How to Refresh a View
Full computing from base data
Efficient if there has been many changes
Incremental computing by applying only the changes to
the view
Better if a small subset has been changed
Uses differential relations which reflect updated data only
© 2020, M.T. Özsu & P. Valduriez 15
Differential Relations
Given relation R and update u
R+ contains tuples inserted by u
R- contains tuples deleted by u
Type of u
insert R- empty
delete R+ empty
modifyR+ (R – R- )
Refreshing a view V is then done by computing
V+ (V – V- )
computing V+ and V- may require accessing base data
© 2020, M.T. Özsu & P. Valduriez 16
Example
EG = SELECT DISTINCT ENAME, RESP
FROM EMP, ASG
WHERE [Link]=[Link]
EG+= (SELECT DISTINCT ENAME, RESP
FROM EMP, ASG+
WHERE [Link]=ASG+.ENO) UNION
(SELECT DISTINCT ENAME, RESP
FROM EMP+, ASG
WHERE EMP+.ENO=[Link]) UNION
(SELECT DISTINCT ENAME, RESP
FROM EMP+, ASG+
WHERE EMP+.ENO=ASG+.ENO)
© 2020, M.T. Özsu & P. Valduriez 17
Techniques for Incremental View
Maintenance
Different techniques depending on:
View expressiveness
Non recursive views: SPJ with duplicate elimination, union and
aggregation
Views with outer join
Recursive views
Most frequent case is non recursive views
Problem: an individual tuple in the view may be derived from
several base tuples
Example: tuple M. Smith, Analyst in EG corresponding to
E2, M. Smith, … in EMP
E2,P1,Analyst,24 and E2,P2,Analyst,6 in ASG
Makes deletion difficult
Solution: Counting
© 2020, M.T. Özsu & P. Valduriez 18
Counting Algorithm
Basic idea
Maintain a count of the number of derivations for each tuple in
the view
Increment (resp. decrement) tuple counts based on insertions
(resp. deletions)
A tuple in the view whose count is zero can be deleted
Algorithm
1. Compute V+ and V- using V, base relations and diff. relations
2. Compute positive in V+ and negative counts in V-
3. Compute V+ (V – V- ), deleting each tuple in V with count=0
Optimal: computes exactly the view tuples that are
inserted or deleted
© 2020, M.T. Özsu & P. Valduriez 19
Exploiting Data Skew
Basic idea
Partition the relations on heavy / light values for join attributes
Threshold depends on data size and user parameter
Maintain the join of different parts using different plans
Most cases done using delta processing (Counting)
Few cases require pre-materialization of auxiliary views
Rebalance the partitions to reflect heavy light changes
Reasons for change:
Much more/less occurrences of a value than before
The heavy/light threshold changes due to change in data size
Update times are amortized to account for occasional rebalancing
© 2020, M.T. Özsu & P. Valduriez 20
Example: Triangle Count
Data model
Relations are functions mapping tuples to multiplicities
Updates also map tuples to multiplicities
Triangle count query
Joins relations R, S and T on common variables
Aggregates away all variables a, b and c
Sums over the product of the multiplicities of matching tuples
Next: Maintenance under single-tuple update to R
Single-tuple update maps to multiplicity
If ( ) then the update is an insert (delete)
© 2020, M.T. Özsu & P. Valduriez 21
Naïve Maintenance for Triangle Count
Compute from scratch
Maintenance time: O(N1.5)
Assuming the input relations have size O(N)
Using existing worst-case optimal join algorithms
No extra space needed
© 2020, M.T. Özsu & P. Valduriez 22
Delta Processing for Triangle Count
Compute the change
Maintenance time: O(N)
Intersect the set of c values paired with b’ in S and with a’ in T
No extra space needed
© 2020, M.T. Özsu & P. Valduriez 23
Materialized View for Triangle Count
Compute the change using materialized views
Maintenance time:
Updates to R: O(1) time to look up in VST
Updates to S and T: O(N) time to maintain VST
Extra O(N2) space needed for the view VST
© 2020, M.T. Özsu & P. Valduriez 24
Data Skew for Triangle Count
For , the triangle count can be maintained with
update time and space.
No algorithm can attain for any .
© 2020, M.T. Özsu & P. Valduriez 25
Heavy/Light Partitioning of Relations
Partition R on a into a light part RL and a heavy part RH
Cardinality bounds
For every value a’:
Also partition S on b and T on c
© 2020, M.T. Özsu & P. Valduriez 26
Maintenance for Skew-Aware Views
For joins of light parts only or heavy parts only
Maintenance using delta processing (Counting)
For joins of a heavy part with a light part
Maintenance using pre-materialized views
Next: Consider one skew-aware view at a time
Single-tuple update to R
© 2020, M.T. Özsu & P. Valduriez 27
Case 1: Light-Light Interaction
Skew-aware views (any partition of R)
Maintenance under update
There are at most c values paired with b’
For each such value c, we check (c,a’) in TL in O(1)
Maintenance time:
© 2020, M.T. Özsu & P. Valduriez 28
Case 2: Heavy-Heavy Interaction
Skew-aware view (any partition of R)
Maintenance under update
There are at most c values paired with a’ in TH
For each such value c, we check (b’,c) in SH in O(1)
Maintenance time:
© 2020, M.T. Özsu & P. Valduriez 29
Case 3: Light-Heavy Interaction
Skew-aware view (any partition of R)
Two possible maintenance plans
1. There are at most c values paired with b’in SL
2. There are at most c values paired with a’ in TH
Maintenance time:
© 2020, M.T. Özsu & P. Valduriez 30
Case 4: Heavy-Light Interaction
Skew-aware view (any partition of R)
Maintenance under update
Materialize auxiliary view
Lookup in the view
Maintenance time
O(1) for the skew-aware view
for the auxiliary view
Size of auxiliary view:
© 2020, M.T. Özsu & P. Valduriez 31
View Self-maintainability
A view is self-maintainable if the base relations need not
be accessed
Not the case for the Counting algorithm
Self-maintainability depends on views’ expressiveness
Most SPJ views are often self-maintainable wrt. deletion and
modification, but not wrt. insertion
Example: a view V is self-maintainable wrt to deletion in R if the
key of R is included in V
© 2020, M.T. Özsu & P. Valduriez 32
Outline
Distributed Data Control
View management
Data security
Semantic integrity control
© 2020, M.T. Özsu & P. Valduriez 33
Data Security
Data protection
Prevents the physical content of data to be understood by
unauthorized users
Uses encryption/decryption techniques (Public key)
Access control
Only authorized users perform operations they are allowed to on
database objects
Discretionary access control (DAC)
Long been provided by DBMS with authorization rules
Multilevel access control (MAC)
Increases security with security levels
© 2020, M.T. Özsu & P. Valduriez 34
Discretionary Access Control
Main actors
Subjects (users, groups of users) who execute operations
Operations (in queries or application programs)
Objects, on which operations are performed
Checking whether a subject may perform an op. on an
object
Authorization= (subject, op. type, object def.)
Defined using GRANT OR REVOKE
Centralized: one single user class (admin.) may grant or revoke
Decentralized, with op. type GRANT
More flexible but recursive revoking process which needs the hierarchy
of grants
© 2020, M.T. Özsu & P. Valduriez 35
Problem with DAC
A malicious user can access unauthorized data through
an authorized user
Example
User A has authorized access to R and S
User B has authorized access to S only
B somehow manages to modify an application program used by
A so it writes R data in S
Then B can read unauthorized data (in S) without violating
authorization rules
Solution: multilevel security based on the famous Bell
and Lapuda model for OS security
© 2020, M.T. Özsu & P. Valduriez 36
Multilevel Access Control
Different security levels (clearances)
Top Secret > Secret > Confidential > Unclassified
Access controlled by 2 rules:
No read up
subject S is allowed to read an object of level L only if level(S) ≥ L
Protect data from unauthorized disclosure, e.g. a subject with secret
clearance cannot read top secret data
No write down:
subject S is allowed to write an object of level L only if level(S) ≤ L
Protect data from unauthorized change, e.g. a subject with top secret
clearance can only write top secret data but not secret data (which could
then contain top secret data)
© 2020, M.T. Özsu & P. Valduriez 37
MAC in Relational DB
A relation can be classified at different levels:
Relation: all tuples have the same clearance
Tuple: every tuple has a clearance
Attribute: every attribute has a clearance
A classified relation is thus multilevel
Appears differently (with different data) to subjects with different
clearances
© 2020, M.T. Özsu & P. Valduriez 38
Example
PROJ*: classified at attribute level
PNO SL1 PNAME SL2 BUDGET SL3 LOC SL4
P1 C Instrumentatio C 150000 C Montreal C
P2 C n C 135000 S New York S
P3 S DB Develop. S 250000 S New York S
CAD/CAM
ROJ* as seen by a subject with confidential clearance
PNO SL1 PNAME SL2 BUDGET SL3 LOC SL4
P1 C Instrumentatio C 150000 C Montreal C
P2 C n C Null C Null C
DB Develop.
© 2020, M.T. Özsu & P. Valduriez 39
Distributed Access Control
Additional problems in a distributed environment
Remote user authentication
Typically using a directory service
Should be replicated at some sites for availability
Management of DAC rules
Problem if users’ group can span multiple sites
Rules stored at some directory based on user groups location
Accessing rules may incur remote queries
Covert channels in MAC
© 2020, M.T. Özsu & P. Valduriez 40
Covert Channels
Indirect means to access unauthorized data
Example
Consider a simple DDB with 2 sites: C (confidential) and S
(secret)
Following the “no write down” rule, an update from a subject with
secret clearance can only be sent to S
Following the “no read up” rule, a read query from the same
subject can be sent to both C and S
But the query may contain secret information (e.g. in a select
predicate), so is a potential covert channel
Solution: replicate part of the DB
So that a site at security level L contains all data that a subject at
level L can access (e.g. S above would replicate the confidential
data so it can entirely process secret queries)
© 2020, M.T. Özsu & P. Valduriez 41
Outline
Distributed Data Control
View management
Data security
Semantic integrity control
© 2020, M.T. Özsu & P. Valduriez 42
Semantic Integrity Control
Maintain database consistency by enforcing a set of
constraints defined on the database.
Structural constraints
Basic semantic properties inherent to a data model e.g., unique
key constraint in relational model
Behavioral constraints
Regulate application behavior, e.g., dependencies in the
relational model
Two components
Integrity constraint specification
Integrity constraint enforcement
© 2020, M.T. Özsu & P. Valduriez 43
Semantic Integrity Control
Procedural
Control embedded in each application program
Declarative
Assertions in predicate calculus
Easy to define constraints
Definition of database consistency clear
But inefficient to check assertions for each update
Limit the search space
Decrease the number of data accesses/assertion
Preventive strategies
Checking at compile time
© 2020, M.T. Özsu & P. Valduriez 44
Constraint Specification Language
Predefined constraints
specify the more common constraints of the relational model
Not-null attribute
ENO NOT NULL IN EMP
Unique key
(ENO, PNO) UNIQUE IN ASG
Foreign key
A key in a relation R is a foreign key if it is a primary key of another
relation S and the existence of any of its values in R is dependent
upon the existence of the same value in S
PNO IN ASG REFERENCES PNO IN PROJ
Functional dependency
ENO IN EMP DETERMINES ENAME
© 2020, M.T. Özsu & P. Valduriez 45
Constraint Specification Language
Precompiled constraints
Express preconditions that must be satisfied by all tuples in a
relation for a given update type
(INSERT, DELETE, MODIFY)
NEW - ranges over new tuples to be inserted
OLD - ranges over old tuples to be deleted
General Form
CHECK ON <relation> [WHEN <update type>]
<qualification>
© 2020, M.T. Özsu & P. Valduriez 46
Constraint Specification Language
Precompiled constraints
Domain constraint
CHECK ON PROJ (BUDGET≥500000 AND BUDGET≤1000000)
Domain constraint on deletion
CHECK ON PROJ WHEN DELETE (BUDGET = 0)
Transition constraint
CHECK ON PROJ ([Link] > [Link] AND
[Link] = [Link])
© 2020, M.T. Özsu & P. Valduriez 47
Constraint Specification Language
General constraints
Constraints that must always be true. Formulae of tuple
relational calculus where all variables are quantified.
General Form
CHECK ON <variable>:<relation>,(<qualification>)
Functional dependency
CHECK ON e1:EMP, e2:EMP
([Link] = [Link] IF [Link] = [Link])
Constraint with aggregate function
CHECK ON g:ASG, j:PROJ
(SUM([Link] WHERE [Link] = [Link]) < 100 IF
[Link] = "CAD/CAM")
© 2020, M.T. Özsu & P. Valduriez 48
Integrity Enforcement
Two methods
Detection
Execute update u: D Du
If Du is inconsistent then
if possible: compensate Du Du’
else
undo Du D
Preventive
Execute u: D Du only if Du will be consistent
Determine valid programs
Determine valid states
© 2020, M.T. Özsu & P. Valduriez 49
Query Modification
Preventive
Add the assertion qualification to the update query
Only applicable to tuple calculus formulae with
universally quantified variables
UPDATE PROJ
SET BUDGET = BUDGET*1.1
WHERE PNAME = "CAD/CAM"
UPDATE PROJ
SET BUDGET = BUDGET*1.1
WHERE PNAME = "CAD/CAM"
AND [Link] ≥ 500000
AND [Link] ≤ 1000000
© 2020, M.T. Özsu & P. Valduriez 50
Compiled Assertions
Triple (R,T,C) where
R relation
T update type (insert, delete, modify)
C assertion on differential relations
Example: Foreign key assertion
g ASG, j PROJ : [Link] = [Link]
Compiled assertions:
(ASG, INSERT, C1), (PROJ, DELETE, C2), (PROJ, MODIFY, C3)
where
C1:NEW ASG+ j PROJ: [Link] = [Link]
C2:g ASG, OLD PROJ- : [Link] ≠ [Link]
C3:g ASG, OLD PROJ- NEW PROJ+:
[Link] ≠[Link] OR [Link] = [Link]
© 2020, M.T. Özsu & P. Valduriez 51
Differential Relations
Given relation R and update u
R+ contains tuples inserted by u
R- contains tuples deleted by u
Type of u
insert R- empty
delete R+ empty
modifyR+ (R – R-)
© 2020, M.T. Özsu & P. Valduriez 52
Differential Relations
Algorithm:
Input: Relation R, update u, compiled assertion Ci
Step 1: Generate differential relations R+ and R–
Step 2: Retrieve the tuples of R+ and R– which do not
satisfy Ci
Step 3: If retrieval is not successful, then the assertion is
valid.
Example :
u is delete on J. Enforcing (EMP, DELETE, C2) :
retrieve all tuples of EMP-
into RESULT
where not(C2)
If RESULT = {}, the assertion is verified
© 2020, M.T. Özsu & P. Valduriez 53
Distributed Integrity Control
Problems:
Definition of constraints
Consideration for fragments
Where to store
Replication
Non-replicated : fragments
Enforcement
Minimize costs
© 2020, M.T. Özsu & P. Valduriez 54
Types of Distributed Assertions
Individual assertions
Single relation, single variable
Domain constraint
Set oriented assertions
Single relation, multi-variable
functional dependency
Multi-relation, multi-variable
foreign key
Assertions involving aggregates
© 2020, M.T. Özsu & P. Valduriez 55
Distributed Integrity Control
Assertion Definition
Similar to the centralized techniques
Transform the assertions to compiled assertions
Assertion Storage
Individual assertions
One relation, only fragments
At each fragment site, check for compatibility
If compatible, store; otherwise reject
If all the sites reject, globally reject
Set-oriented assertions
Involves joins (between fragments or relations)
May be necessary to perform joins to check for compatibility
Store if compatible
© 2020, M.T. Özsu & P. Valduriez 56
Distributed Integrity Control
Assertion Enforcement
Where to enforce each assertion depends on
Type of assertion
Type of update and where update is issued
Individual Assertions
If update = insert
Enforce at the site where the update is issued
If update = qualified
Send the assertions to all the sites involved
Execute the qualification to obtain R+ and R-
Each site enforces its own assertion
Set-oriented Assertions
Single relation
Similar to individual assertions with qualified updates
Multi-relation
Move data to perform joins; then send the result to query master site
© 2020, M.T. Özsu & P. Valduriez 57
Conclusion
Solutions initially designed for centralized systems have
been significantly extended for distributed systems
Materialized views and group-based discretionary access control
Semantic integrity control has received less attention
and is generally not well supported by distributed DBMS
products
Full data control is more complex and costly in
distributed systems
Definition and storage of the rules (site selection)
Design of enforcement algorithms which minimize
communication costs
© 2020, M.T. Özsu & P. Valduriez 58