Advanced SQL Interview Questions and Answers for Experienced Developers
Interview preparation · Technical guide
Advanced SQL Interview Questions and Answers
Questions, explanations, and illustrative examples for interview preparation.
Examples are independent teaching snippets and may require application types, imports, packages, schema, and configuration. Framework behavior is version-dependent. Corrections address identified issues; the complete source code collection has not been compiled or integration-tested.
1. How would you design a database schema for a multi-tenant SaaS application?
Answer: There are three main approaches: Shared Database/Shared Schema, Shared Database/Separate Schemas, and Separate Databases.
Recommended Approach: Shared Database/Separate Schemas
-- Tenant management
CREATE TABLE tenants (
tenant_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
tenant_name VARCHAR(255) NOT NULL,
schema_name VARCHAR(63) NOT NULL UNIQUE,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
is_active BOOLEAN DEFAULT TRUE
);
-- Tenant-specific schema creation function
CREATE OR REPLACE FUNCTION create_tenant_schema(tenant_schema_name VARCHAR)
RETURNS VOID AS $$
BEGIN
EXECUTE 'CREATE SCHEMA IF NOT EXISTS ' || quote_ident(tenant_schema_name);
-- Create tenant-specific tables
EXECUTE 'CREATE TABLE ' || quote_ident(tenant_schema_name) || '.users (
user_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
email VARCHAR(255) UNIQUE NOT NULL,
name VARCHAR(255) NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)';
EXECUTE 'CREATE TABLE ' || quote_ident(tenant_schema_name) || '.projects (
project_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
name VARCHAR(255) NOT NULL,
description TEXT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)';
END;
$$ LANGUAGE plpgsql;
-- Row-level security for shared tables
CREATE TABLE shared_users (
user_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
tenant_id UUID REFERENCES tenants(tenant_id),
email VARCHAR(255) NOT NULL,
name VARCHAR(255) NOT NULL
);
-- Enable RLS
ALTER TABLE shared_users ENABLE ROW LEVEL SECURITY;
-- Policy to restrict access by tenant
CREATE POLICY tenant_isolation_policy ON shared_users
FOR ALL USING (tenant_id = current_setting('app.current_tenant_id')::UUID);
2. What are the different normalization forms and when would you denormalize?
Answer: Normalization forms (1NF, 2NF, 3NF, BCNF, 4NF, 5NF) reduce redundancy and anomalies.
When to Denormalize: - Read-heavy workloads - Complex analytical queries - Performance optimization - Data warehousing
-- Normalized schema (3NF)
CREATE TABLE orders (
order_id INT PRIMARY KEY,
customer_id INT,
order_date DATE,
total_amount DECIMAL(10,2)
);
CREATE TABLE order_items (
order_item_id INT PRIMARY KEY,
order_id INT REFERENCES orders(order_id),
product_id INT,
quantity INT,
unit_price DECIMAL(10,2)
);
CREATE TABLE products (
product_id INT PRIMARY KEY,
name VARCHAR(255),
category_id INT
);
CREATE TABLE categories (
category_id INT PRIMARY KEY,
name VARCHAR(255)
);
-- Denormalized for reporting (violates 3NF)
CREATE TABLE order_summary (
order_id INT PRIMARY KEY,
customer_id INT,
customer_name VARCHAR(255), -- Denormalized
order_date DATE,
total_amount DECIMAL(10,2),
product_count INT, -- Denormalized
category_names TEXT -- Denormalized (comma-separated)
);
-- Materialized view for complex analytics
CREATE MATERIALIZED VIEW sales_by_category AS
SELECT
c.name as category_name,
COUNT(o.order_id) as order_count,
SUM(oi.quantity * oi.unit_price) as total_revenue,
AVG(oi.unit_price) as avg_price
FROM orders o
JOIN order_items oi ON o.order_id = oi.order_id
JOIN products p ON oi.product_id = p.product_id
JOIN categories c ON p.category_id = c.category_id
GROUP BY c.category_id, c.name
WITH DATA;
-- Refresh materialized view
REFRESH MATERIALIZED VIEW sales_by_category;
3. How do you handle database versioning and schema migrations?
Answer: Use migration tools, version control, and rollback strategies.
Example using Alembic (Python/SQLAlchemy)
migrations/versions/001_create_users_table.py
"""Create users table
Revision ID: 001 Revises: Create Date: 2024-01-01 10:00:00.000000
""" from alembic import op import sqlalchemy as sa
revision identifiers
revision = '001' down_revision = None branch_labels = None depends_on = None
def upgrade(): op.create_table('users', sa.Column('id', sa.Integer(), nullable=False), sa.Column('email', sa.String(length=255), nullable=False), sa.Column('name', sa.String(length=255), nullable=False), sa.Column('created_at', sa.DateTime(), nullable=False), sa.PrimaryKeyConstraint('id'), sa.UniqueConstraint('email') )
# Create index for performance
op.create_index('idx_users_email', 'users', ['email'])
def downgrade(): op.drop_index('idx_users_email', table_name='users') op.drop_table('users')
-- Database versioning table
CREATE TABLE schema_migrations (
version VARCHAR(255) PRIMARY KEY,
applied_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
checksum VARCHAR(64),
description TEXT
);
-- Migration tracking function
CREATE OR REPLACE FUNCTION apply_migration(
migration_version VARCHAR,
migration_sql TEXT,
migration_checksum VARCHAR
) RETURNS BOOLEAN AS $$
BEGIN
-- Check if migration already applied
IF EXISTS (SELECT 1 FROM schema_migrations WHERE version = migration_version) THEN
RETURN FALSE;
END IF;
-- Apply migration
EXECUTE migration_sql;
-- Record migration
INSERT INTO schema_migrations (version, checksum)
VALUES (migration_version, migration_checksum);
RETURN TRUE;
END;
$$ LANGUAGE plpgsql;
4. What's the difference between OLTP and OLAP databases?
Answer: OLTP (Online Transaction Processing) vs OLAP (Online Analytical Processing)
OLTP Database Design:
-- Optimized for transactions
CREATE TABLE orders (
order_id BIGSERIAL PRIMARY KEY,
customer_id INT NOT NULL,
order_date TIMESTAMP NOT NULL,
status VARCHAR(20) NOT NULL,
total_amount DECIMAL(10,2) NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- Indexes for fast lookups
CREATE INDEX idx_orders_customer_id ON orders(customer_id);
CREATE INDEX idx_orders_status ON orders(status);
CREATE INDEX idx_orders_date ON orders(order_date);
-- Optimized for INSERT/UPDATE/DELETE
INSERT INTO orders (customer_id, order_date, status, total_amount)
VALUES (123, CURRENT_TIMESTAMP, 'pending', 99.99);
OLAP Database Design (Star Schema):
-- Fact table
CREATE TABLE sales_fact (
sale_id BIGSERIAL PRIMARY KEY,
date_id INT REFERENCES date_dim(date_id),
customer_id INT REFERENCES customer_dim(customer_id),
product_id INT REFERENCES product_dim(product_id),
store_id INT REFERENCES store_dim(store_id),
quantity INT,
unit_price DECIMAL(10,2),
total_amount DECIMAL(10,2)
);
-- Dimension tables
CREATE TABLE date_dim (
date_id INT PRIMARY KEY,
full_date DATE,
year INT,
quarter INT,
month INT,
day_of_week INT,
is_holiday BOOLEAN
);
CREATE TABLE customer_dim (
customer_id INT PRIMARY KEY,
customer_name VARCHAR(255),
age_group VARCHAR(20),
income_level VARCHAR(20),
location VARCHAR(255)
);
-- Optimized for complex queries
CREATE INDEX idx_sales_date ON sales_fact(date_id);
CREATE INDEX idx_sales_customer ON sales_fact(customer_id);
CREATE INDEX idx_sales_product ON sales_fact(product_id);
-- Materialized view for common aggregations
CREATE MATERIALIZED VIEW monthly_sales_summary AS
SELECT
d.year,
d.month,
p.category,
SUM(sf.quantity) as total_quantity,
SUM(sf.total_amount) as total_revenue
FROM sales_fact sf
JOIN date_dim d ON sf.date_id = d.date_id
JOIN product_dim p ON sf.product_id = p.product_id
GROUP BY d.year, d.month, p.category
WITH DATA;
5. How would you design a database for a social media platform?
Answer: Focus on scalability, performance, and handling massive data volumes.
-- User management with partitioning
CREATE TABLE users (
user_id BIGSERIAL PRIMARY KEY,
username VARCHAR(50) UNIQUE NOT NULL,
email VARCHAR(255) UNIQUE NOT NULL,
password_hash VARCHAR(255) NOT NULL,
profile_picture_url VARCHAR(500),
bio TEXT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
is_verified BOOLEAN DEFAULT FALSE,
follower_count INT DEFAULT 0,
following_count INT DEFAULT 0
) PARTITION BY HASH (user_id);
-- Create partitions
CREATE TABLE users_p0 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 0);
CREATE TABLE users_p1 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 1);
CREATE TABLE users_p2 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 2);
CREATE TABLE users_p3 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 3);
-- Posts with time-based partitioning
CREATE TABLE posts (
post_id BIGSERIAL PRIMARY KEY,
user_id BIGINT NOT NULL,
content TEXT,
media_urls TEXT[], -- Array of media URLs
location_lat DECIMAL(10,8),
location_lng DECIMAL(11,8),
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
like_count INT DEFAULT 0,
comment_count INT DEFAULT 0,
share_count INT DEFAULT 0
) PARTITION BY RANGE (created_at);
-- Monthly partitions
CREATE TABLE posts_2024_01 PARTITION OF posts
FOR VALUES FROM ('2024-01-01') TO ('2024-02-01');
CREATE TABLE posts_2024_02 PARTITION OF posts
FOR VALUES FROM ('2024-02-01') TO ('2024-03-01');
-- Relationships (followers/following)
CREATE TABLE user_relationships (
follower_id BIGINT NOT NULL,
following_id BIGINT NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (follower_id, following_id),
FOREIGN KEY (follower_id) REFERENCES users(user_id),
FOREIGN KEY (following_id) REFERENCES users(user_id)
);
-- Likes with composite primary key
CREATE TABLE post_likes (
post_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (post_id, user_id),
FOREIGN KEY (post_id) REFERENCES posts(post_id),
FOREIGN KEY (user_id) REFERENCES users(user_id)
);
-- Comments with hierarchical structure
CREATE TABLE comments (
comment_id BIGSERIAL PRIMARY KEY,
post_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
parent_comment_id BIGINT, -- For nested comments
content TEXT NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
like_count INT DEFAULT 0,
FOREIGN KEY (post_id) REFERENCES posts(post_id),
FOREIGN KEY (user_id) REFERENCES users(user_id),
FOREIGN KEY (parent_comment_id) REFERENCES comments(comment_id)
);
-- Indexes for performance
CREATE INDEX idx_posts_user_created ON posts(user_id, created_at DESC);
CREATE INDEX idx_posts_created_at ON posts(created_at DESC);
CREATE INDEX idx_user_relationships_follower ON user_relationships(follower_id);
CREATE INDEX idx_user_relationships_following ON user_relationships(following_id);
CREATE INDEX idx_post_likes_post ON post_likes(post_id);
CREATE INDEX idx_comments_post_created ON comments(post_id, created_at DESC);
-- Feed generation view
CREATE VIEW user_feed AS
SELECT
p.post_id,
p.user_id,
u.username,
p.content,
p.media_urls,
p.created_at,
p.like_count,
p.comment_count
FROM posts p
JOIN users u ON p.user_id = u.user_id
JOIN user_relationships ur ON p.user_id = ur.following_id
WHERE ur.follower_id = current_setting('app.current_user_id')::BIGINT
ORDER BY p.created_at DESC;
6. What are the considerations for designing a distributed database system?
Answer: Focus on consistency, availability, partition tolerance (CAP theorem), and data distribution strategies.
-- Sharding strategy example
-- Global sequence for cross-shard IDs
CREATE SEQUENCE global_id_seq;
-- Shard metadata
CREATE TABLE shard_metadata (
shard_id INT PRIMARY KEY,
shard_name VARCHAR(50) NOT NULL,
host VARCHAR(255) NOT NULL,
port INT NOT NULL,
database_name VARCHAR(50) NOT NULL,
key_range_start BIGINT,
key_range_end BIGINT,
is_active BOOLEAN DEFAULT TRUE
);
-- Distributed users table (sharded by user_id)
CREATE TABLE users (
user_id BIGINT PRIMARY KEY DEFAULT nextval('global_id_seq'),
username VARCHAR(50) UNIQUE NOT NULL,
email VARCHAR(255) UNIQUE NOT NULL,
shard_id INT NOT NULL, -- For routing
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- Shard routing function
CREATE OR REPLACE FUNCTION get_shard_for_user(user_id BIGINT)
RETURNS INT AS $$
BEGIN
RETURN (user_id % 4) + 1; -- 4 shards
END;
$$ LANGUAGE plpgsql;
-- Cross-shard query coordination
CREATE OR REPLACE FUNCTION distributed_user_search(search_term VARCHAR)
RETURNS TABLE (
user_id BIGINT,
username VARCHAR(50),
email VARCHAR(255),
shard_id INT
) AS $$
DECLARE
shard_record RECORD;
BEGIN
FOR shard_record IN SELECT * FROM shard_metadata WHERE is_active = TRUE
LOOP
-- Execute query on each shard
RETURN QUERY EXECUTE format(
'SELECT user_id, username, email, shard_id
FROM users@%s
WHERE username ILIKE %L OR email ILIKE %L',
shard_record.shard_name,
'%' || search_term || '%',
'%' || search_term || '%'
);
END LOOP;
END;
$$ LANGUAGE plpgsql;
Consistency patterns
class DistributedDatabase: def init(self): self.shards = {} # shard_id -> connection self.consistency_level = 'eventual' # or 'strong'
def write_with_consistency(self, data, consistency_level='quorum'):
"""Write with specified consistency level"""
if consistency_level == 'strong':
# Write to all replicas synchronously
for shard in self.shards.values():
shard.write(data)
elif consistency_level == 'quorum':
# Write to majority of replicas
replicas = list(self.shards.values())
majority = len(replicas) // 2 + 1
for i in range(majority):
replicas[i].write(data)
else: # eventual
# Write to primary, replicate asynchronously
primary = self.get_primary_shard()
primary.write(data)
def read_with_consistency(self, key, consistency_level='eventual'):
"""Read with specified consistency level"""
if consistency_level == 'strong':
# Read from all replicas and compare
results = []
for shard in self.shards.values():
results.append(shard.read(key))
return self.resolve_conflicts(results)
else:
# Read from nearest replica
return self.get_nearest_shard().read(key)
7. How do you implement soft deletes vs hard deletes?
Answer: Soft deletes preserve data integrity and enable recovery, while hard deletes permanently remove data.
-- Soft delete implementation
CREATE TABLE products (
product_id BIGSERIAL PRIMARY KEY,
name VARCHAR(255) NOT NULL,
price DECIMAL(10,2) NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
deleted_at TIMESTAMP NULL, -- Soft delete timestamp
deleted_by BIGINT NULL, -- Who deleted it
is_deleted BOOLEAN DEFAULT FALSE -- Boolean flag
);
-- Index for soft delete queries
CREATE INDEX idx_products_deleted_at ON products(deleted_at) WHERE deleted_at IS NULL;
CREATE INDEX idx_products_is_deleted ON products(is_deleted) WHERE is_deleted = FALSE;
-- Soft delete function
CREATE OR REPLACE FUNCTION soft_delete_product(
p_product_id BIGINT,
p_deleted_by BIGINT
) RETURNS BOOLEAN AS $$
BEGIN
UPDATE products
SET
deleted_at = CURRENT_TIMESTAMP,
deleted_by = p_deleted_by,
is_deleted = TRUE,
updated_at = CURRENT_TIMESTAMP
WHERE product_id = p_product_id AND is_deleted = FALSE;
RETURN FOUND;
END;
$$ LANGUAGE plpgsql;
-- Restore function
CREATE OR REPLACE FUNCTION restore_product(p_product_id BIGINT)
RETURNS BOOLEAN AS $$
BEGIN
UPDATE products
SET
deleted_at = NULL,
deleted_by = NULL,
is_deleted = FALSE,
updated_at = CURRENT_TIMESTAMP
WHERE product_id = p_product_id AND is_deleted = TRUE;
RETURN FOUND;
END;
$$ LANGUAGE plpgsql;
-- View for active products only
CREATE VIEW active_products AS
SELECT * FROM products WHERE is_deleted = FALSE;
-- Hard delete function (use with caution)
CREATE OR REPLACE FUNCTION hard_delete_product(p_product_id BIGINT)
RETURNS BOOLEAN AS $$
BEGIN
DELETE FROM products WHERE product_id = p_product_id;
RETURN FOUND;
END;
$$ LANGUAGE plpgsql;
-- Cleanup old soft-deleted records (archival)
CREATE OR REPLACE FUNCTION archive_old_deleted_products(months_old INT DEFAULT 12)
RETURNS INT AS $$
DECLARE
archived_count INT;
BEGIN
-- Move to archive table
INSERT INTO products_archive
SELECT * FROM products
WHERE deleted_at < CURRENT_TIMESTAMP - INTERVAL '1 month' * months_old
AND is_deleted = TRUE;
GET DIAGNOSTICS archived_count = ROW_COUNT;
-- Delete from main table
DELETE FROM products
WHERE deleted_at < CURRENT_TIMESTAMP - INTERVAL '1 month' * months_old
AND is_deleted = TRUE;
RETURN archived_count;
END;
$$ LANGUAGE plpgsql;
8. What's the difference between a clustered and non-clustered index?
Answer: Clustered indexes determine physical storage order, while non-clustered indexes are separate structures.
-- Clustered index (determines physical order)
CREATE TABLE orders (
order_id BIGSERIAL PRIMARY KEY, -- Automatically clustered
customer_id BIGINT NOT NULL,
order_date DATE NOT NULL,
total_amount DECIMAL(10,2) NOT NULL,
status VARCHAR(20) NOT NULL
);
-- Create clustered index on order_date (reorganizes table)
CREATE INDEX idx_orders_date_clustered ON orders(order_date);
-- Note: In PostgreSQL, you'd use CLUSTER command:
-- CLUSTER orders USING idx_orders_date_clustered;
-- Non-clustered indexes (separate structures)
CREATE INDEX idx_orders_customer ON orders(customer_id);
CREATE INDEX idx_orders_status ON orders(status);
CREATE INDEX idx_orders_date_status ON orders(order_date, status);
-- Composite non-clustered index
CREATE INDEX idx_orders_customer_date ON orders(customer_id, order_date DESC);
-- Partial index (non-clustered)
CREATE INDEX idx_orders_active ON orders(customer_id, order_date)
WHERE status = 'active';
-- Covering index (includes all columns needed for query)
CREATE INDEX idx_orders_covering ON orders(customer_id, order_date, total_amount, status);
-- Unique non-clustered index
CREATE UNIQUE INDEX idx_orders_unique_reference ON orders(customer_id, order_date);
-- Index with included columns (SQL Server style - PostgreSQL uses covering indexes)
-- CREATE INDEX idx_orders_included ON orders(customer_id, order_date)
-- INCLUDE (total_amount, status);
-- Performance comparison queries
-- Query using clustered index (fastest)
EXPLAIN ANALYZE SELECT * FROM orders WHERE order_date = '2024-01-15';
-- Query using non-clustered index
EXPLAIN ANALYZE SELECT * FROM orders WHERE customer_id = 123;
-- Query using composite index
EXPLAIN ANALYZE SELECT * FROM orders
WHERE customer_id = 123 AND order_date >= '2024-01-01';
-- Query using covering index (no table lookup needed)
EXPLAIN ANALYZE SELECT customer_id, order_date, total_amount
FROM orders WHERE customer_id = 123;
-- Index usage statistics
SELECT
schemaname,
tablename,
indexname,
idx_scan as index_scans,
idx_tup_read as tuples_read,
idx_tup_fetch as tuples_fetched
FROM pg_stat_user_indexes
WHERE tablename = 'orders'
ORDER BY idx_scan DESC;
9. How would you design a database for an e-commerce platform?
Answer: Focus on performance, scalability, and handling complex business logic.
-- Product catalog with categories
CREATE TABLE categories (
category_id BIGSERIAL PRIMARY KEY,
name VARCHAR(255) NOT NULL,
parent_category_id BIGINT REFERENCES categories(category_id),
slug VARCHAR(255) UNIQUE NOT NULL,
is_active BOOLEAN DEFAULT TRUE,
sort_order INT DEFAULT 0
);
CREATE TABLE products (
product_id BIGSERIAL PRIMARY KEY,
sku VARCHAR(100) UNIQUE NOT NULL,
name VARCHAR(255) NOT NULL,
description TEXT,
short_description VARCHAR(500),
category_id BIGINT REFERENCES categories(category_id),
brand_id BIGINT,
price DECIMAL(10,2) NOT NULL,
compare_price DECIMAL(10,2),
cost_price DECIMAL(10,2),
weight DECIMAL(8,3),
dimensions JSONB, -- {length, width, height}
is_active BOOLEAN DEFAULT TRUE,
is_featured BOOLEAN DEFAULT FALSE,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- Product variants (sizes, colors, etc.)
CREATE TABLE product_variants (
variant_id BIGSERIAL PRIMARY KEY,
product_id BIGINT REFERENCES products(product_id),
sku VARCHAR(100) UNIQUE NOT NULL,
variant_name VARCHAR(255), -- e.g., "Red, Large"
price_adjustment DECIMAL(10,2) DEFAULT 0,
stock_quantity INT DEFAULT 0,
is_active BOOLEAN DEFAULT TRUE
);
-- Inventory management
CREATE TABLE inventory (
inventory_id BIGSERIAL PRIMARY KEY,
product_id BIGINT REFERENCES products(product_id),
variant_id BIGINT REFERENCES product_variants(variant_id),
warehouse_id BIGINT,
quantity_available INT DEFAULT 0,
quantity_reserved INT DEFAULT 0,
reorder_point INT DEFAULT 0,
last_updated TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- Customer management
CREATE TABLE customers (
customer_id BIGSERIAL PRIMARY KEY,
email VARCHAR(255) UNIQUE NOT NULL,
password_hash VARCHAR(255),
first_name VARCHAR(100),
last_name VARCHAR(100),
phone VARCHAR(20),
date_of_birth DATE,
is_active BOOLEAN DEFAULT TRUE,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- Addresses
CREATE TABLE addresses (
address_id BIGSERIAL PRIMARY KEY,
customer_id BIGINT REFERENCES customers(customer_id),
address_type VARCHAR(20) DEFAULT 'shipping', -- shipping, billing
first_name VARCHAR(100),
last_name VARCHAR(100),
company VARCHAR(255),
address_line1 VARCHAR(255),
address_line2 VARCHAR(255),
city VARCHAR(100),
state VARCHAR(100),
postal_code VARCHAR(20),
country VARCHAR(100),
phone VARCHAR(20),
is_default BOOLEAN DEFAULT FALSE
);
-- Orders
CREATE TABLE orders (
order_id BIGSERIAL PRIMARY KEY,
order_number VARCHAR(50) UNIQUE NOT NULL,
customer_id BIGINT REFERENCES customers(customer_id),
status VARCHAR(50) NOT NULL DEFAULT 'pending', -- pending, confirmed, shipped, delivered, cancelled
subtotal DECIMAL(10,2) NOT NULL,
tax_amount DECIMAL(10,2) DEFAULT 0,
shipping_amount DECIMAL(10,2) DEFAULT 0,
discount_amount DECIMAL(10,2) DEFAULT 0,
total_amount DECIMAL(10,2) NOT NULL,
shipping_address_id BIGINT REFERENCES addresses(address_id),
billing_address_id BIGINT REFERENCES addresses(address_id),
payment_method VARCHAR(50),
payment_status VARCHAR(50) DEFAULT 'pending',
notes TEXT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- Order items
CREATE TABLE order_items (
order_item_id BIGSERIAL PRIMARY KEY,
order_id BIGINT REFERENCES orders(order_id),
product_id BIGINT REFERENCES products(product_id),
variant_id BIGINT REFERENCES product_variants(variant_id),
sku VARCHAR(100) NOT NULL,
product_name VARCHAR(255) NOT NULL,
quantity INT NOT NULL,
unit_price DECIMAL(10,2) NOT NULL,
total_price DECIMAL(10,2) NOT NULL,
tax_amount DECIMAL(10,2) DEFAULT 0
);
-- Shopping cart
CREATE TABLE cart_items (
cart_item_id BIGSERIAL PRIMARY KEY,
customer_id BIGINT REFERENCES customers(customer_id),
session_id VARCHAR(255), -- For guest users
product_id BIGINT REFERENCES products(product_id),
variant_id BIGINT REFERENCES product_variants(variant_id),
quantity INT NOT NULL DEFAULT 1,
added_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- Performance indexes
CREATE INDEX idx_products_category ON products(category_id) WHERE is_active = TRUE;
CREATE INDEX idx_products_price ON products(price) WHERE is_active = TRUE;
CREATE INDEX idx_products_featured ON products(is_featured) WHERE is_featured = TRUE;
CREATE INDEX idx_orders_customer_status ON orders(customer_id, status);
CREATE INDEX idx_orders_created_at ON orders(created_at DESC);
CREATE INDEX idx_order_items_order ON order_items(order_id);
CREATE INDEX idx_inventory_product ON inventory(product_id, variant_id);
CREATE INDEX idx_cart_items_customer ON cart_items(customer_id);
CREATE INDEX idx_cart_items_session ON cart_items(session_id);
-- Full-text search index
CREATE INDEX idx_products_search ON products USING gin(
to_tsvector('english', name || ' ' || COALESCE(description, ''))
);
-- Materialized view for product analytics
CREATE MATERIALIZED VIEW product_analytics AS
SELECT
p.product_id,
p.name,
p.category_id,
COUNT(oi.order_item_id) as total_orders,
SUM(oi.quantity) as total_quantity_sold,
SUM(oi.total_price) as total_revenue,
AVG(oi.unit_price) as avg_selling_price
FROM products p
LEFT JOIN order_items oi ON p.product_id = oi.product_id
LEFT JOIN orders o ON oi.order_id = o.order_id
WHERE o.status IN ('confirmed', 'shipped', 'delivered')
GROUP BY p.product_id, p.name, p.category_id
WITH DATA;
10. What are the best practices for database partitioning?
Answer: Partitioning improves performance and manageability for large tables.
-- Range partitioning by date
CREATE TABLE sales (
sale_id BIGSERIAL,
customer_id BIGINT,
product_id BIGINT,
sale_date DATE NOT NULL,
amount DECIMAL(10,2),
region VARCHAR(50)
) PARTITION BY RANGE (sale_date);
-- Create partitions for each month
CREATE TABLE sales_2024_01 PARTITION OF sales
FOR VALUES FROM ('2024-01-01') TO ('2024-02-01');
CREATE TABLE sales_2024_02 PARTITION OF sales
FOR VALUES FROM ('2024-02-01') TO ('2024-03-01');
CREATE TABLE sales_2024_03 PARTITION OF sales
FOR VALUES FROM ('2024-03-01') TO ('2024-04-01');
-- Hash partitioning for even distribution
CREATE TABLE users (
user_id BIGSERIAL,
username VARCHAR(50),
email VARCHAR(255),
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
) PARTITION BY HASH (user_id);
-- Create hash partitions
CREATE TABLE users_p0 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 0);
CREATE TABLE users_p1 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 1);
CREATE TABLE users_p2 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 2);
CREATE TABLE users_p3 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 3);
-- List partitioning by region
CREATE TABLE orders (
order_id BIGSERIAL,
customer_id BIGINT,
region VARCHAR(50) NOT NULL,
order_date DATE,
total_amount DECIMAL(10,2)
) PARTITION BY LIST (region);
CREATE TABLE orders_us PARTITION OF orders FOR VALUES IN ('US');
CREATE TABLE orders_eu PARTITION OF orders FOR VALUES IN ('EU', 'UK');
CREATE TABLE orders_asia PARTITION OF orders FOR VALUES IN ('JP', 'CN', 'IN');
-- Composite partitioning (range + hash)
CREATE TABLE events (
event_id BIGSERIAL,
event_date DATE NOT NULL,
user_id BIGINT,
event_type VARCHAR(50),
event_data JSONB
) PARTITION BY RANGE (event_date);
-- Sub-partition by hash
CREATE TABLE events_2024_01 PARTITION OF events
FOR VALUES FROM ('2024-01-01') TO ('2024-02-01')
PARTITION BY HASH (user_id);
CREATE TABLE events_2024_01_p0 PARTITION OF events_2024_01
FOR VALUES WITH (modulus 4, remainder 0);
CREATE TABLE events_2024_01_p1 PARTITION OF events_2024_01
FOR VALUES WITH (modulus 4, remainder 1);
-- Partition management functions
CREATE OR REPLACE FUNCTION create_monthly_partition(
table_name TEXT,
partition_date DATE
) RETURNS VOID AS $$
DECLARE
partition_name TEXT;
start_date DATE;
end_date DATE;
BEGIN
partition_name := table_name || '_' || to_char(partition_date, 'YYYY_MM');
start_date := date_trunc('month', partition_date);
end_date := start_date + INTERVAL '1 month';
EXECUTE format(
'CREATE TABLE %I PARTITION OF %I FOR VALUES FROM (%L) TO (%L)',
partition_name, table_name, start_date, end_date
);
RAISE NOTICE 'Created partition %', partition_name;
END;
$$ LANGUAGE plpgsql;
-- Auto-create partitions
CREATE OR REPLACE FUNCTION auto_create_partitions() RETURNS TRIGGER AS $$
BEGIN
-- Create partition for current month if it doesn't exist
PERFORM create_monthly_partition('sales', CURRENT_DATE);
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER trigger_auto_create_partitions
BEFORE INSERT ON sales
FOR EACH ROW
EXECUTE FUNCTION auto_create_partitions();
-- Partition maintenance
CREATE OR REPLACE FUNCTION drop_old_partitions(
table_name TEXT,
months_to_keep INT DEFAULT 12
) RETURNS INT AS $$
DECLARE
partition_record RECORD;
dropped_count INT := 0;
cutoff_date DATE;
BEGIN
cutoff_date := CURRENT_DATE - INTERVAL '1 month' * months_to_keep;
FOR partition_record IN
SELECT tablename
FROM pg_tables
WHERE tablename LIKE table_name || '_%'
AND tablename ~ '^\d{4}_\d{2}$'
LOOP
-- Extract date from partition name and check if old
IF to_date(split_part(partition_record.tablename, '_', 2) || '_' ||
split_part(partition_record.tablename, '_', 3), 'YYYY_MM') < cutoff_date THEN
EXECUTE 'DROP TABLE ' || partition_record.tablename;
dropped_count := dropped_count + 1;
END IF;
END LOOP;
RETURN dropped_count;
END;
$$ LANGUAGE plpgsql;
-- Partition statistics
CREATE VIEW partition_stats AS
SELECT
schemaname,
tablename,
attname,
n_distinct,
correlation
FROM pg_stats
WHERE tablename LIKE '%_2024_%'
ORDER BY tablename, attname;
-- Partition performance monitoring
SELECT
schemaname,
tablename,
seq_scan,
seq_tup_read,
idx_scan,
idx_tup_fetch,
n_tup_ins,
n_tup_upd,
n_tup_del
FROM pg_stat_user_tables
WHERE tablename LIKE 'sales_%'
ORDER BY tablename;
11. How do you identify and fix slow-running queries?
Identification Methods:
- Dynamic Management Views (DMVs): Query sys.dm_exec_query_stats for execution statistics
- Extended Events: Capture query performance data
- SQL Server Profiler: Real-time monitoring
- Query Store: Built-in performance monitoring (SQL Server 2016+)
Example - Identifying Slow Queries:
-- Find top 10 slowest queries
SELECT TOP 10
qs.total_elapsed_time / qs.execution_count AS avg_elapsed_time,
qs.execution_count,
qs.total_logical_reads / qs.execution_count AS avg_logical_reads,
SUBSTRING(qt.text, (qs.statement_start_offset/2)+1,
((CASE qs.statement_end_offset
WHEN -1 THEN DATALENGTH(qt.text)
ELSE qs.statement_end_offset
END - qs.statement_start_offset)/2) + 1) AS query_text
FROM sys.dm_exec_query_stats qs
CROSS APPLY sys.dm_exec_sql_text(qs.sql_handle) qt
ORDER BY avg_elapsed_time DESC;
Fixing Strategies: 1. Add missing indexes 2. Rewrite queries to use more efficient patterns 3. Update statistics 4. Optimize table structure
12. What are the different types of SQL Server execution plans?
Types of Execution Plans:
- Estimated Execution Plan: Shows what SQL Server thinks the query will do
- Actual Execution Plan: Shows what actually happened during execution
- Live Query Statistics: Real-time execution monitoring
Key Plan Elements:
-- Generate estimated execution plan
SET SHOWPLAN_ALL ON;
GO
SELECT * FROM Orders o
INNER JOIN Customers c ON o.CustomerID = c.CustomerID
WHERE o.OrderDate >= '2023-01-01';
GO
SET SHOWPLAN_ALL OFF;
-- Generate actual execution plan
SET STATISTICS IO ON;
SET STATISTICS TIME ON;
SELECT * FROM Orders o
INNER JOIN Customers c ON o.CustomerID = c.CustomerID
WHERE o.OrderDate >= '2023-01-01';
Common Plan Operators: - Table Scan: Reads entire table - Index Scan: Reads entire index - Index Seek: Direct lookup using index - Hash Join: For large datasets - Nested Loop Join: For small datasets - Sort: Memory or disk-based sorting
13. How do you optimize a query that's performing table scans?
Table Scan Optimization Strategies:
- Add Clustered Index:
-- Create clustered index on frequently queried column
CREATE CLUSTERED INDEX IX_Orders_OrderDate
ON Orders(OrderDate);
-- Composite clustered index
CREATE CLUSTERED INDEX IX_Orders_Composite
ON Orders(OrderDate, CustomerID);
- Add Non-Clustered Indexes:
-- Covering index for specific query
CREATE NONCLUSTERED INDEX IX_Orders_Covering
ON Orders(OrderDate, CustomerID)
INCLUDE (OrderAmount, OrderStatus);
- Query Rewriting:
-- Before (causes table scan)
SELECT * FROM Orders WHERE YEAR(OrderDate) = 2023;
-- After (uses index)
SELECT * FROM Orders
WHERE OrderDate >= '2023-01-01'
AND OrderDate < '2024-01-01';
- Partitioning:
-- Partition by date range
CREATE PARTITION FUNCTION PF_OrderDate (datetime)
AS RANGE RIGHT FOR VALUES ('2023-01-01', '2024-01-01', '2025-01-01');
14. What's the difference between covering and non-covering indexes?
Covering Index: - Contains all columns needed by the query - Eliminates key lookups - Faster execution
Non-Covering Index: - Missing some required columns - Requires additional lookups - Slower execution
Example:
-- Query requiring multiple columns
SELECT OrderID, CustomerID, OrderDate, OrderAmount
FROM Orders
WHERE OrderDate >= '2023-01-01';
-- Non-covering index (requires key lookup)
CREATE NONCLUSTERED INDEX IX_Orders_OrderDate
ON Orders(OrderDate);
-- Covering index (no key lookup needed)
CREATE NONCLUSTERED INDEX IX_Orders_Covering
ON Orders(OrderDate)
INCLUDE (OrderID, CustomerID, OrderAmount);
Performance Comparison:
-- Check if index is covering
SELECT
i.name AS IndexName,
i.type_desc,
CASE
WHEN ic.is_included_column = 1 THEN 'INCLUDED'
ELSE 'KEY'
END AS ColumnType,
c.name AS ColumnName
FROM sys.indexes i
INNER JOIN sys.index_columns ic ON i.object_id = ic.object_id AND i.index_id = ic.index_id
INNER JOIN sys.columns c ON ic.object_id = c.object_id AND ic.column_id = c.column_id
WHERE i.object_id = OBJECT_ID('Orders');
15. How do you handle deadlocks in SQL Server?
Deadlock Prevention Strategies:
- Consistent Access Order:
-- Always access tables in same order
BEGIN TRANSACTION;
UPDATE Customers SET LastOrderDate = GETDATE() WHERE CustomerID = 1;
UPDATE Orders SET OrderStatus = 'Processed' WHERE CustomerID = 1;
COMMIT;
- Use NOLOCK Hint (with caution):
-- For read-only operations where data consistency isn't critical
SELECT * FROM Orders WITH (NOLOCK) WHERE OrderDate >= '2023-01-01';
- Deadlock Monitoring:
-- Monitor deadlocks
SELECT
deadlock_graph.value('(/event/@timestamp)[1]', 'datetime2') AS DeadlockTime,
deadlock_graph.value('(/event/data/value)[1]', 'varchar(max)') AS DeadlockGraph
FROM (
SELECT CAST(target_data AS xml) AS deadlock_graph
FROM sys.dm_xe_session_targets st
JOIN sys.dm_xe_sessions s ON s.address = st.event_session_address
WHERE s.name = 'system_health'
AND st.target_name = 'ring_buffer'
) AS DeadlockData;
- Deadlock Priority:
-- Set deadlock priority
SET DEADLOCK_PRIORITY HIGH; -- Current session has higher priority
-- or
SET DEADLOCK_PRIORITY LOW; -- Current session has lower priority
16. What are the performance implications of different join types?
Join Types and Performance:
- INNER JOIN:
-- Most efficient when both tables have indexes on join columns
SELECT o.OrderID, c.CustomerName
FROM Orders o
INNER JOIN Customers c ON o.CustomerID = c.CustomerID
WHERE o.OrderDate >= '2023-01-01';
- LEFT/RIGHT JOIN:
-- Can be slower due to NULL handling
SELECT c.CustomerName, o.OrderID
FROM Customers c
LEFT JOIN Orders o ON c.CustomerID = o.CustomerID
WHERE o.OrderDate >= '2023-01-01' OR o.OrderID IS NULL;
- CROSS JOIN:
-- Cartesian product - very expensive
SELECT p.ProductName, c.CategoryName
FROM Products p
CROSS JOIN Categories c; -- Avoid unless necessary
- Hash Join Optimization:
-- Force hash join for large datasets
SELECT o.OrderID, c.CustomerName
FROM Orders o
INNER HASH JOIN Customers c ON o.CustomerID = c.CustomerID
WHERE o.OrderDate >= '2023-01-01';
Performance Tips: - Use appropriate indexes on join columns - Consider query hints for large datasets - Monitor execution plans for join type selection
17. How do you optimize stored procedures?
Stored Procedure Optimization Techniques:
- Parameter Sniffing Issues:
-- Use local variables to avoid parameter sniffing
CREATE PROCEDURE GetOrdersByDate
@StartDate datetime,
@EndDate datetime
AS
BEGIN
DECLARE @LocalStartDate datetime = @StartDate;
DECLARE @LocalEndDate datetime = @EndDate;
SELECT OrderID, CustomerID, OrderDate
FROM Orders
WHERE OrderDate BETWEEN @LocalStartDate AND @LocalEndDate;
END;
- Recompilation Strategies:
-- Force recompilation when needed
CREATE PROCEDURE GetOrdersByDate
@StartDate datetime,
@EndDate datetime
WITH RECOMPILE
AS
BEGIN
SELECT OrderID, CustomerID, OrderDate
FROM Orders
WHERE OrderDate BETWEEN @StartDate AND @EndDate;
END;
- SET NOCOUNT ON:
CREATE PROCEDURE GetOrdersByDate
@StartDate datetime,
@EndDate datetime
AS
BEGIN
SET NOCOUNT ON; -- Reduces network traffic
SELECT OrderID, CustomerID, OrderDate
FROM Orders
WHERE OrderDate BETWEEN @StartDate AND @EndDate;
END;
- Error Handling:
CREATE PROCEDURE GetOrdersByDate
@StartDate datetime,
@EndDate datetime
AS
BEGIN
SET NOCOUNT ON;
BEGIN TRY
SELECT OrderID, CustomerID, OrderDate
FROM Orders
WHERE OrderDate BETWEEN @StartDate AND @EndDate;
END TRY
BEGIN CATCH
SELECT
ERROR_NUMBER() AS ErrorNumber,
ERROR_MESSAGE() AS ErrorMessage;
END CATCH;
END;
18. What's the impact of parameter sniffing and how do you handle it?
Parameter Sniffing Impact: - SQL Server creates execution plan based on first parameter values - Subsequent calls with different parameters may use suboptimal plan - Can cause performance degradation
Handling Strategies:
- Use Local Variables:
CREATE PROCEDURE GetOrdersByDate
@StartDate datetime,
@EndDate datetime
AS
BEGIN
DECLARE @LocalStartDate datetime = @StartDate;
DECLARE @LocalEndDate datetime = @EndDate;
SELECT OrderID, CustomerID, OrderDate
FROM Orders
WHERE OrderDate BETWEEN @LocalStartDate AND @LocalEndDate;
END;
- OPTION (RECOMPILE):
CREATE PROCEDURE GetOrdersByDate
@StartDate datetime,
@EndDate datetime
AS
BEGIN
SELECT OrderID, CustomerID, OrderDate
FROM Orders
WHERE OrderDate BETWEEN @StartDate AND @EndDate
OPTION (RECOMPILE);
END;
- OPTION (OPTIMIZE FOR UNKNOWN):
CREATE PROCEDURE GetOrdersByDate
@StartDate datetime,
@EndDate datetime
AS
BEGIN
SELECT OrderID, CustomerID, OrderDate
FROM Orders
WHERE OrderDate BETWEEN @StartDate AND @EndDate
OPTION (OPTIMIZE FOR UNKNOWN);
END;
- Plan Guides:
-- Create plan guide to force specific plan
EXEC sp_create_plan_guide
@name = N'GetOrdersByDate_PlanGuide',
@stmt = N'SELECT OrderID, CustomerID, OrderDate FROM Orders WHERE OrderDate BETWEEN @StartDate AND @EndDate',
@type = N'OBJECT',
@module_or_batch = N'GetOrdersByDate',
@params = NULL,
@hints = N'OPTION (RECOMPILE)';
19. How do you implement query result caching?
Caching Implementation Strategies:
- Application-Level Caching:
// C# example using MemoryCache
public class OrderService
{
private readonly IMemoryCache _cache;
public async Task<List<Order>> GetOrdersAsync(DateTime startDate, DateTime endDate)
{
string cacheKey = $"Orders_{startDate:yyyyMMdd}_{endDate:yyyyMMdd}";
if (_cache.TryGetValue(cacheKey, out List<Order> cachedOrders))
{
return cachedOrders;
}
var orders = await GetOrdersFromDatabaseAsync(startDate, endDate);
var cacheOptions = new MemoryCacheEntryOptions()
.SetSlidingExpiration(TimeSpan.FromMinutes(30))
.SetAbsoluteExpiration(TimeSpan.FromHours(2));
_cache.Set(cacheKey, orders, cacheOptions);
return orders;
}
}
- Redis Caching:
// C# example using Redis
public class OrderService
{
private readonly IDistributedCache _cache;
public async Task<List<Order>> GetOrdersAsync(DateTime startDate, DateTime endDate)
{
string cacheKey = $"Orders_{startDate:yyyyMMdd}_{endDate:yyyyMMdd}";
var cachedData = await _cache.GetStringAsync(cacheKey);
if (!string.IsNullOrEmpty(cachedData))
{
return JsonSerializer.Deserialize<List<Order>>(cachedData);
}
var orders = await GetOrdersFromDatabaseAsync(startDate, endDate);
var options = new DistributedCacheEntryOptions()
.SetSlidingExpiration(TimeSpan.FromMinutes(30));
await _cache.SetStringAsync(cacheKey, JsonSerializer.Serialize(orders), options);
return orders;
}
}
- SQL Server Query Store:
-- Enable Query Store
ALTER DATABASE YourDatabase
SET QUERY_STORE = ON
(
OPERATION_MODE = READ_WRITE,
CLEANUP_POLICY = (STALE_QUERY_THRESHOLD_DAYS = 30),
DATA_FLUSH_INTERVAL_SECONDS = 3000,
MAX_STORAGE_SIZE_MB = 1000,
INTERVAL_LENGTH_MINUTES = 10
);
-- Query cached plans
SELECT
qsq.query_id,
qsq.query_hash,
qsp.plan_id,
qsp.query_plan,
qsrs.avg_duration,
qsrs.avg_logical_io_reads
FROM sys.query_store_query qsq
INNER JOIN sys.query_store_plan qsp ON qsq.query_id = qsp.query_id
INNER JOIN sys.query_store_runtime_stats qsrs ON qsp.plan_id = qsrs.plan_id
ORDER BY qsrs.avg_duration DESC;
20. What are the best practices for database connection pooling?
Connection Pooling Best Practices:
- Connection String Configuration:
// C# connection string with pooling settings
string connectionString =
"Server=myserver;Database=mydb;Integrated Security=true;" +
"Max Pool Size=100;" +
"Min Pool Size=10;" +
"Pooling=true;" +
"Connection Lifetime=300;" +
"Connection Reset=false;";
- Using Statement Best Practices:
// Proper connection disposal
public async Task<List<Order>> GetOrdersAsync()
{
using (var connection = new SqlConnection(connectionString))
{
await connection.OpenAsync();
using (var command = new SqlCommand("SELECT * FROM Orders", connection))
using (var reader = await command.ExecuteReaderAsync())
{
var orders = new List<Order>();
while (await reader.ReadAsync())
{
orders.Add(new Order
{
OrderID = reader.GetInt32("OrderID"),
CustomerID = reader.GetInt32("CustomerID"),
OrderDate = reader.GetDateTime("OrderDate")
});
}
return orders;
}
}
}
- Dependency Injection Configuration:
// ASP.NET Core service configuration
public void ConfigureServices(IServiceCollection services)
{
services.AddDbContext<ApplicationDbContext>(options =>
options.UseSqlServer(Configuration.GetConnectionString("DefaultConnection"),
sqlServerOptionsAction: sqlOptions =>
{
sqlOptions.EnableRetryOnFailure(
maxRetryCount: 3,
maxRetryDelay: TimeSpan.FromSeconds(30),
errorNumbersToAdd: null);
}));
}
- Monitoring Connection Pool:
-- Monitor connection pool usage
SELECT
DB_NAME(dbid) AS DatabaseName,
COUNT(dbid) AS NumberOfConnections,
loginame AS LoginName
FROM sys.sysprocesses
WHERE dbid > 0
GROUP BY dbid, loginame
ORDER BY NumberOfConnections DESC;
-- Check connection pool settings
SELECT
name,
value,
value_in_use
FROM sys.configurations
WHERE name LIKE '%pool%';
Key Best Practices: - Set appropriate pool sizes based on application load - Always dispose connections properly - Use connection string parameters for tuning - Monitor pool usage and performance - Implement retry logic for transient failures - Use async/await for better scalability
21. How do you write a recursive CTE for hierarchical data?
Answer: Recursive CTEs (Common Table Expressions) are used to query hierarchical data like organizational charts, file systems, or category trees. They consist of an anchor member (base case) and a recursive member.
Example - Employee Hierarchy:
WITH RECURSIVE EmployeeHierarchy AS (
-- Anchor member: Find the root employee (CEO)
SELECT
employee_id,
name,
manager_id,
0 AS level,
CAST(name AS VARCHAR(1000)) AS hierarchy_path
FROM employees
WHERE manager_id IS NULL
UNION ALL
-- Recursive member: Find all subordinates
SELECT
e.employee_id,
e.name,
e.manager_id,
eh.level + 1,
CAST(eh.hierarchy_path + ' > ' + e.name AS VARCHAR(1000))
FROM employees e
INNER JOIN EmployeeHierarchy eh ON e.manager_id = eh.employee_id
WHERE eh.level < 10 -- Prevent infinite recursion
)
SELECT
level,
hierarchy_path,
employee_id,
name
FROM EmployeeHierarchy
ORDER BY hierarchy_path;
Example - File System Structure:
WITH RECURSIVE FileTree AS (
-- Anchor: Root directories
SELECT
file_id,
file_name,
parent_id,
0 AS depth,
CAST(file_name AS VARCHAR(500)) AS path
FROM files
WHERE parent_id IS NULL
UNION ALL
-- Recursive: Subdirectories and files
SELECT
f.file_id,
f.file_name,
f.parent_id,
ft.depth + 1,
CAST(ft.path + '/' + f.file_name AS VARCHAR(500))
FROM files f
INNER JOIN FileTree ft ON f.parent_id = ft.file_id
)
SELECT
REPEAT(' ', depth) || file_name AS display_name,
path,
depth
FROM FileTree
ORDER BY path;
22. How do you implement pivot and unpivot operations?
Answer: Pivot transforms rows into columns, while unpivot does the opposite. These are essential for data transformation and reporting.
Pivot Example:
-- Pivot: Transform sales data from rows to columns
SELECT
product_name,
SUM(CASE WHEN year = 2021 THEN sales_amount END) AS sales_2021,
SUM(CASE WHEN year = 2022 THEN sales_amount END) AS sales_2022,
SUM(CASE WHEN year = 2023 THEN sales_amount END) AS sales_2023,
SUM(CASE WHEN year = 2024 THEN sales_amount END) AS sales_2024
FROM sales_data
GROUP BY product_name;
-- Using PIVOT operator (SQL Server)
SELECT product_name, [2021], [2022], [2023], [2024]
FROM (
SELECT product_name, year, sales_amount
FROM sales_data
) AS source
PIVOT (
SUM(sales_amount)
FOR year IN ([2021], [2022], [2023], [2024])
) AS pvt;
Unpivot Example:
-- Unpivot: Transform columns back to rows
SELECT
product_name,
year,
sales_amount
FROM (
SELECT
product_name,
sales_2021,
sales_2022,
sales_2023,
sales_2024
FROM sales_pivot_table
) AS source
UNPIVOT (
sales_amount FOR year IN (sales_2021, sales_2022, sales_2023, sales_2024)
) AS unpvt;
-- Using CROSS JOIN for multiple columns
SELECT
product_name,
year,
sales_amount,
quantity
FROM sales_pivot_table
CROSS JOIN (
SELECT 2021 AS year, 'sales_2021' AS sales_col, 'qty_2021' AS qty_col
UNION ALL SELECT 2022, 'sales_2022', 'qty_2022'
UNION ALL SELECT 2023, 'sales_2023', 'qty_2023'
UNION ALL SELECT 2024, 'sales_2024', 'qty_2024'
) AS years;
23. How do you handle dynamic SQL safely?
Answer: Dynamic SQL should be handled carefully to prevent SQL injection and ensure proper error handling.
Safe Dynamic SQL Implementation:
-- Using parameterized queries (SQL Server)
CREATE PROCEDURE GetEmployeeData
@TableName NVARCHAR(128),
@Department NVARCHAR(50) = NULL
AS
BEGIN
SET NOCOUNT ON;
-- Validate table name to prevent injection
IF @TableName NOT IN ('employees', 'contractors', 'managers')
BEGIN
RAISERROR('Invalid table name', 16, 1);
RETURN;
END
DECLARE @SQL NVARCHAR(MAX);
DECLARE @Params NVARCHAR(500) = N'@Dept NVARCHAR(50)';
SET @SQL = N'
SELECT employee_id, name, department, salary
FROM ' + QUOTENAME(@TableName) + '
WHERE (@Dept IS NULL OR department = @Dept)
ORDER BY name';
EXEC sp_executesql @SQL, @Params, @Dept = @Department;
END;
PostgreSQL Example:
-- Using EXECUTE with format() for safe dynamic SQL
CREATE OR REPLACE FUNCTION get_employee_data(
table_name TEXT,
department TEXT DEFAULT NULL
) RETURNS TABLE (
employee_id INTEGER,
name TEXT,
department TEXT,
salary NUMERIC
) AS $$
DECLARE
sql_query TEXT;
BEGIN
-- Validate table name
IF table_name NOT IN ('employees', 'contractors', 'managers') THEN
RAISE EXCEPTION 'Invalid table name: %', table_name;
END IF;
sql_query := format(
'SELECT employee_id, name, department, salary
FROM %I
WHERE ($1 IS NULL OR department = $1)
ORDER BY name',
table_name
);
RETURN QUERY EXECUTE sql_query USING department;
END;
$$ LANGUAGE plpgsql;
24. How do you implement row-level security?
Answer: Row-level security (RLS) restricts data access based on user context, ensuring users only see data they're authorized to access.
SQL Server RLS Implementation:
-- Enable RLS on table
ALTER TABLE employees ENABLE ROW LEVEL SECURITY;
-- Create security policy
CREATE SECURITY POLICY EmployeeSecurityPolicy
ADD FILTER PREDICATE dbo.fn_securitypredicate(employee_id)
ON dbo.employees;
-- Security predicate function
CREATE FUNCTION dbo.fn_securitypredicate(@employee_id INT)
RETURNS TABLE
WITH SCHEMABINDING
AS
RETURN SELECT 1 AS fn_securitypredicate_result
WHERE @employee_id IN (
-- User can see their own records
SELECT employee_id FROM employees WHERE user_id = SYSTEM_USER
UNION
-- Managers can see their team
SELECT e.employee_id
FROM employees e
INNER JOIN employees m ON e.manager_id = m.employee_id
WHERE m.user_id = SYSTEM_USER
UNION
-- HR can see all records
SELECT employee_id FROM employees
WHERE EXISTS (
SELECT 1 FROM user_roles
WHERE user_id = SYSTEM_USER AND role = 'HR'
)
);
PostgreSQL RLS Example:
-- Enable RLS
ALTER TABLE employees ENABLE ROW LEVEL SECURITY;
-- Create policies
CREATE POLICY employee_access_policy ON employees
FOR ALL
USING (
-- User can see their own record
employee_id = current_setting('app.current_user_id')::INT
OR
-- Manager can see team members
manager_id = current_setting('app.current_user_id')::INT
OR
-- HR can see all
EXISTS (
SELECT 1 FROM user_roles
WHERE user_id = current_setting('app.current_user_id')::INT
AND role = 'HR'
)
);
-- Set user context
SELECT set_config('app.current_user_id', '123', false);
25. How do you write a query to find gaps in sequential data?
Answer: Finding gaps in sequential data is common for identifying missing records, unused IDs, or discontinuous sequences.
Method 1: Using Window Functions
-- Find gaps in employee IDs
WITH numbered_rows AS (
SELECT
employee_id,
ROW_NUMBER() OVER (ORDER BY employee_id) AS expected_id
FROM employees
)
SELECT
expected_id AS missing_id,
LAG(employee_id) OVER (ORDER BY expected_id) AS prev_id,
LEAD(employee_id) OVER (ORDER BY expected_id) AS next_id
FROM numbered_rows
WHERE employee_id != expected_id;
Method 2: Using Self-Join
-- Find gaps in sequence
SELECT
t1.sequence_value + 1 AS gap_start,
MIN(t2.sequence_value) - 1 AS gap_end
FROM sequence_table t1
LEFT JOIN sequence_table t2 ON t1.sequence_value < t2.sequence_value
WHERE t2.sequence_value IS NULL
OR t2.sequence_value > t1.sequence_value + 1
GROUP BY t1.sequence_value
HAVING MIN(t2.sequence_value) - 1 >= t1.sequence_value + 1;
Method 3: Using Recursive CTE
-- Find all missing values in a range
WITH RECURSIVE missing_values AS (
SELECT MIN(sequence_value) AS value
FROM sequence_table
UNION ALL
SELECT mv.value + 1
FROM missing_values mv
WHERE mv.value < (SELECT MAX(sequence_value) FROM sequence_table)
)
SELECT mv.value AS missing_value
FROM missing_values mv
WHERE NOT EXISTS (
SELECT 1 FROM sequence_table st
WHERE st.sequence_value = mv.value
);
26. How do you implement a running total calculation?
Answer: Running totals are calculated using window functions to show cumulative values over ordered data.
Basic Running Total:
-- Running total of sales by date
SELECT
order_date,
sales_amount,
SUM(sales_amount) OVER (
ORDER BY order_date
ROWS UNBOUNDED PRECEDING
) AS running_total,
SUM(sales_amount) OVER (
ORDER BY order_date
RANGE UNBOUNDED PRECEDING
) AS running_total_same_date
FROM sales_orders
ORDER BY order_date;
Running Total by Category:
-- Running total by product category
SELECT
order_date,
product_category,
sales_amount,
SUM(sales_amount) OVER (
PARTITION BY product_category
ORDER BY order_date
ROWS UNBOUNDED PRECEDING
) AS category_running_total,
SUM(sales_amount) OVER (
ORDER BY order_date
ROWS UNBOUNDED PRECEDING
) AS overall_running_total
FROM sales_orders
ORDER BY product_category, order_date;
Moving Average with Running Total:
-- 30-day moving average with running total
SELECT
order_date,
sales_amount,
SUM(sales_amount) OVER (
ORDER BY order_date
ROWS UNBOUNDED PRECEDING
) AS running_total,
AVG(sales_amount) OVER (
ORDER BY order_date
ROWS BETWEEN 29 PRECEDING AND CURRENT ROW
) AS moving_30day_avg
FROM sales_orders
ORDER BY order_date;
27. How do you handle multiple date ranges in a single query?
Answer: Multiple date ranges can be handled using various techniques like date dimension tables, CTEs, or complex WHERE clauses.
Using Date Dimension Table:
-- Create date dimension
CREATE TABLE date_dimension (
date_key DATE PRIMARY KEY,
year INT,
quarter INT,
month INT,
week INT,
day_of_week INT,
is_weekend BOOLEAN,
is_holiday BOOLEAN
);
-- Query with multiple date ranges
SELECT
product_id,
SUM(CASE
WHEN d.date_key BETWEEN '2024-01-01' AND '2024-03-31'
THEN sales_amount END) AS q1_sales,
SUM(CASE
WHEN d.date_key BETWEEN '2024-04-01' AND '2024-06-30'
THEN sales_amount END) AS q2_sales,
SUM(CASE
WHEN d.date_key BETWEEN '2024-07-01' AND '2024-09-30'
THEN sales_amount END) AS q3_sales,
SUM(CASE
WHEN d.date_key BETWEEN '2024-10-01' AND '2024-12-31'
THEN sales_amount END) AS q4_sales
FROM sales s
JOIN date_dimension d ON s.order_date = d.date_key
WHERE d.year = 2024
GROUP BY product_id;
Using CTE for Date Ranges:
WITH date_ranges AS (
SELECT
'Q1' AS quarter,
'2024-01-01'::DATE AS start_date,
'2024-03-31'::DATE AS end_date
UNION ALL
SELECT 'Q2', '2024-04-01'::DATE, '2024-06-30'::DATE
UNION ALL
SELECT 'Q3', '2024-07-01'::DATE, '2024-09-30'::DATE
UNION ALL
SELECT 'Q4', '2024-10-01'::DATE, '2024-12-31'::DATE
)
SELECT
p.product_name,
dr.quarter,
SUM(s.sales_amount) AS quarterly_sales
FROM products p
CROSS JOIN date_ranges dr
LEFT JOIN sales s ON p.product_id = s.product_id
AND s.order_date BETWEEN dr.start_date AND dr.end_date
GROUP BY p.product_name, dr.quarter, dr.start_date
ORDER BY p.product_name, dr.start_date;
28. How do you implement data archiving strategies?
Answer: Data archiving involves moving old data to separate storage while maintaining accessibility and compliance.
Partitioned Archiving Strategy:
-- Create partitioned table for archiving
CREATE TABLE sales_archive (
order_id INT,
order_date DATE,
customer_id INT,
amount DECIMAL(10,2)
) PARTITION BY RANGE (order_date);
-- Create partitions for different years
CREATE TABLE sales_archive_2020 PARTITION OF sales_archive
FOR VALUES FROM ('2020-01-01') TO ('2021-01-01');
CREATE TABLE sales_archive_2021 PARTITION OF sales_archive
FOR VALUES FROM ('2021-01-01') TO ('2022-01-01');
-- Archive procedure
CREATE OR REPLACE PROCEDURE archive_old_sales(archive_date DATE)
LANGUAGE plpgsql
AS $$
BEGIN
-- Move old data to archive
INSERT INTO sales_archive
SELECT * FROM sales
WHERE order_date < archive_date;
-- Delete from main table
DELETE FROM sales
WHERE order_date < archive_date;
-- Update statistics
ANALYZE sales;
ANALYZE sales_archive;
END;
$$;
Temporal Table Approach:
-- Create temporal table with system time
CREATE TABLE employee_history (
employee_id INT,
name VARCHAR(100),
department VARCHAR(50),
salary DECIMAL(10,2),
valid_from TIMESTAMP GENERATED ALWAYS AS ROW START,
valid_to TIMESTAMP GENERATED ALWAYS AS ROW END,
PERIOD FOR SYSTEM_TIME (valid_from, valid_to)
) WITH (SYSTEM_VERSIONING = ON);
-- Archive old versions
CREATE TABLE employee_archive (
employee_id INT,
name VARCHAR(100),
department VARCHAR(50),
salary DECIMAL(10,2),
valid_from TIMESTAMP,
valid_to TIMESTAMP
);
-- Archive procedure
CREATE PROCEDURE archive_employee_history(@cutoff_date DATE)
AS
BEGIN
INSERT INTO employee_archive
SELECT * FROM employee_history
FOR SYSTEM_TIME ALL
WHERE valid_to < @cutoff_date;
-- Clean up old history
DELETE FROM employee_history
WHERE valid_to < @cutoff_date;
END;
29. How do you write a query to find duplicate records?
Answer: Finding duplicates requires identifying records with identical values in specified columns.
Method 1: Using GROUP BY
-- Find duplicates based on email
SELECT
email,
COUNT(*) AS duplicate_count,
STRING_AGG(CAST(user_id AS VARCHAR), ', ') AS user_ids
FROM users
GROUP BY email
HAVING COUNT(*) > 1;
Method 2: Using Window Functions
-- Find duplicates with row details
WITH duplicate_check AS (
SELECT
*,
ROW_NUMBER() OVER (
PARTITION BY email, first_name, last_name
ORDER BY created_date
) AS rn,
COUNT(*) OVER (
PARTITION BY email, first_name, last_name
) AS duplicate_count
FROM users
)
SELECT
user_id,
email,
first_name,
last_name,
created_date,
duplicate_count
FROM duplicate_check
WHERE duplicate_count > 1
ORDER BY email, created_date;
Method 3: Self-Join Approach
-- Find exact duplicates
SELECT DISTINCT
u1.user_id AS user_id_1,
u2.user_id AS user_id_2,
u1.email,
u1.first_name,
u1.last_name
FROM users u1
INNER JOIN users u2 ON
u1.email = u2.email
AND u1.first_name = u2.first_name
AND u1.last_name = u2.last_name
AND u1.user_id < u2.user_id;
30. How do you implement data validation constraints?
Answer: Data validation constraints ensure data integrity at the database level.
Check Constraints:
-- Table with comprehensive constraints
CREATE TABLE employees (
employee_id INT PRIMARY KEY IDENTITY(1,1),
first_name VARCHAR(50) NOT NULL,
last_name VARCHAR(50) NOT NULL,
email VARCHAR(100) NOT NULL UNIQUE,
phone VARCHAR(20),
salary DECIMAL(10,2) NOT NULL,
hire_date DATE NOT NULL,
department VARCHAR(50),
-- Check constraints
CONSTRAINT chk_salary_positive CHECK (salary > 0),
CONSTRAINT chk_hire_date_valid CHECK (hire_date <= GETDATE()),
CONSTRAINT chk_email_format CHECK (email LIKE '%_@_%._%'),
CONSTRAINT chk_phone_format CHECK (
phone IS NULL OR
phone LIKE '[0-9][0-9][0-9]-[0-9][0-9][0-9]-[0-9][0-9][0-9][0-9]'
),
CONSTRAINT chk_department_valid CHECK (
department IN ('IT', 'HR', 'Finance', 'Marketing', 'Sales')
)
);
Custom Validation Functions:
-- PostgreSQL custom validation function
CREATE OR REPLACE FUNCTION validate_credit_card(card_number TEXT)
RETURNS BOOLEAN AS $$
DECLARE
sum INTEGER := 0;
digit INTEGER;
i INTEGER;
BEGIN
-- Luhn algorithm implementation
FOR i IN REVERSE LENGTH(card_number)..1 LOOP
digit := CAST(SUBSTRING(card_number FROM i FOR 1) AS INTEGER);
IF (LENGTH(card_number) - i) % 2 = 0 THEN
digit := digit * 2;
IF digit > 9 THEN
digit := digit - 9;
END IF;
END IF;
sum := sum + digit;
END LOOP;
RETURN sum % 10 = 0;
END;
$$ LANGUAGE plpgsql;
-- Use in table constraint
CREATE TABLE customers (
customer_id SERIAL PRIMARY KEY,
name VARCHAR(100) NOT NULL,
credit_card VARCHAR(20),
CONSTRAINT chk_valid_credit_card
CHECK (credit_card IS NULL OR validate_credit_card(credit_card))
);
31. How do you implement database backup and recovery strategies?
Answer: Backup and recovery strategies ensure data protection and business continuity.
SQL Server Backup Strategy:
-- Full backup procedure
CREATE PROCEDURE PerformFullBackup
@DatabaseName NVARCHAR(128),
@BackupPath NVARCHAR(500)
AS
BEGIN
DECLARE @BackupFileName NVARCHAR(500);
DECLARE @BackupName NVARCHAR(500);
SET @BackupFileName = @BackupPath + '\' + @DatabaseName + '_' +
CONVERT(VARCHAR(8), GETDATE(), 112) + '_' +
REPLACE(CONVERT(VARCHAR(8), GETDATE(), 108), ':', '') + '.bak';
SET @BackupName = @DatabaseName + ' Full Backup ' +
CONVERT(VARCHAR(20), GETDATE(), 120);
BACKUP DATABASE @DatabaseName
TO DISK = @BackupFileName
WITH
NAME = @BackupName,
DESCRIPTION = 'Full database backup',
COMPRESSION,
CHECKSUM,
STATS = 10;
END;
-- Differential backup procedure
CREATE PROCEDURE PerformDifferentialBackup
@DatabaseName NVARCHAR(128),
@BackupPath NVARCHAR(500)
AS
BEGIN
DECLARE @BackupFileName NVARCHAR(500);
SET @BackupFileName = @BackupPath + '\' + @DatabaseName + '_diff_' +
CONVERT(VARCHAR(8), GETDATE(), 112) + '.bak';
BACKUP DATABASE @DatabaseName
TO DISK = @BackupFileName
WITH
DIFFERENTIAL,
COMPRESSION,
CHECKSUM;
END;
PostgreSQL Backup Strategy:
-- Backup function using pg_dump
CREATE OR REPLACE FUNCTION perform_backup(
db_name TEXT,
backup_path TEXT
) RETURNS TEXT AS $$
DECLARE
backup_file TEXT;
cmd TEXT;
result TEXT;
BEGIN
backup_file := backup_path || '/' || db_name || '_' ||
TO_CHAR(CURRENT_TIMESTAMP, 'YYYYMMDD_HH24MISS') || '.sql';
cmd := 'pg_dump -h localhost -U postgres -d ' || db_name ||
' -f ' || backup_file || ' --verbose --no-password';
-- Execute backup command
SELECT pg_exec(cmd) INTO result;
RETURN 'Backup completed: ' || backup_file;
END;
$$ LANGUAGE plpgsql;
32. What are the different types of database replication?
Answer: Database replication types include master-slave, multi-master, and various synchronization strategies.
Master-Slave Replication Setup:
-- Master server configuration (MySQL)
-- my.cnf on master
[mysqld]
server-id = 1
log-bin = mysql-bin
binlog_format = ROW
sync_binlog = 1
-- Slave server configuration
[mysqld]
server-id = 2
relay-log = mysql-relay-bin
read_only = 1
-- Setup replication
-- On master
CREATE USER 'repl'@'%' IDENTIFIED BY 'password';
GRANT REPLICATION SLAVE ON *.* TO 'repl'@'%';
FLUSH PRIVILEGES;
-- On slave
CHANGE MASTER TO
MASTER_HOST = 'master_host',
MASTER_USER = 'repl',
MASTER_PASSWORD = 'password',
MASTER_LOG_FILE = 'mysql-bin.000001',
MASTER_LOG_POS = 154;
START SLAVE;
PostgreSQL Logical Replication:
-- Publisher setup
CREATE PUBLICATION sales_pub FOR TABLE sales, customers;
-- Subscriber setup
CREATE SUBSCRIPTION sales_sub
CONNECTION 'host=subscriber_host port=5432 dbname=target_db user=repl password=password'
PUBLICATION sales_pub;
-- Monitor replication
SELECT * FROM pg_stat_subscription;
SELECT * FROM pg_publication_tables;
33. How do you handle data consistency in distributed systems?
Answer: Data consistency in distributed systems requires careful coordination and conflict resolution strategies.
Two-Phase Commit Implementation:
-- Coordinator table for 2PC
CREATE TABLE transaction_coordinator (
transaction_id UUID PRIMARY KEY,
status VARCHAR(20) DEFAULT 'PREPARING',
participants JSONB,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- Participant table
CREATE TABLE transaction_participants (
transaction_id UUID,
participant_id VARCHAR(100),
status VARCHAR(20),
resource_id VARCHAR(100),
PRIMARY KEY (transaction_id, participant_id)
);
-- 2PC Procedure
CREATE OR REPLACE PROCEDURE execute_2pc_transaction(
p_transaction_id UUID,
p_operations JSONB
)
LANGUAGE plpgsql
AS $$
DECLARE
participant RECORD;
all_prepared BOOLEAN := TRUE;
BEGIN
-- Phase 1: Prepare
INSERT INTO transaction_coordinator (transaction_id, participants)
VALUES (p_transaction_id, p_operations);
-- Send prepare to all participants
FOR participant IN SELECT * FROM jsonb_array_elements(p_operations)
LOOP
-- Call prepare on each participant
PERFORM prepare_transaction(participant);
IF participant.status != 'PREPARED' THEN
all_prepared := FALSE;
END IF;
END LOOP;
-- Phase 2: Commit or Abort
IF all_prepared THEN
-- Commit all participants
UPDATE transaction_coordinator
SET status = 'COMMITTED', updated_at = CURRENT_TIMESTAMP
WHERE transaction_id = p_transaction_id;
-- Send commit to all participants
FOR participant IN SELECT * FROM jsonb_array_elements(p_operations)
LOOP
PERFORM commit_transaction(participant);
END LOOP;
ELSE
-- Abort all participants
UPDATE transaction_coordinator
SET status = 'ABORTED', updated_at = CURRENT_TIMESTAMP
WHERE transaction_id = p_transaction_id;
FOR participant IN SELECT * FROM jsonb_array_elements(p_operations)
LOOP
PERFORM abort_transaction(participant);
END LOOP;
END IF;
END;
$$;
34. How do you implement data encryption at rest and in transit?
Answer: Data encryption protects sensitive information both when stored and during transmission.
SQL Server Encryption:
-- Create master key
CREATE MASTER KEY ENCRYPTION BY PASSWORD = 'StrongPassword123!';
-- Create certificate
CREATE CERTIFICATE MyCert
WITH SUBJECT = 'Database Encryption Certificate';
-- Create database encryption key
CREATE DATABASE ENCRYPTION KEY
WITH ALGORITHM = AES_256
ENCRYPTION BY SERVER CERTIFICATE MyCert;
-- Enable transparent data encryption
ALTER DATABASE MyDatabase
SET ENCRYPTION ON;
-- Column-level encryption
CREATE SYMMETRIC KEY CreditCardKey
WITH ALGORITHM = AES_256
ENCRYPTION BY CERTIFICATE MyCert;
-- Encrypt sensitive data
UPDATE customers
SET credit_card_encrypted = ENCRYPTBYKEY(KEY_GUID('CreditCardKey'), credit_card)
WHERE credit_card IS NOT NULL;
-- Decrypt data
SELECT
customer_id,
name,
CAST(DECRYPTBYKEY(credit_card_encrypted) AS VARCHAR(20)) AS credit_card
FROM customers;
PostgreSQL Encryption:
-- Enable pgcrypto extension
CREATE EXTENSION IF NOT EXISTS pgcrypto;
-- Encrypt function
CREATE OR REPLACE FUNCTION encrypt_sensitive_data(
data TEXT,
encryption_key TEXT
) RETURNS TEXT AS $$
BEGIN
RETURN encode(encrypt_iv(
data::bytea,
encryption_key::bytea,
decode('1234567890123456', 'hex'),
'aes-cbc'
), 'base64');
END;
$$ LANGUAGE plpgsql;
-- Decrypt function
CREATE OR REPLACE FUNCTION decrypt_sensitive_data(
encrypted_data TEXT,
encryption_key TEXT
) RETURNS TEXT AS $$
BEGIN
RETURN convert_from(decrypt_iv(
decode(encrypted_data, 'base64'),
encryption_key::bytea,
decode('1234567890123456', 'hex'),
'aes-cbc'
), 'utf8');
END;
$$ LANGUAGE plpgsql;
-- Use in table
CREATE TABLE secure_customers (
customer_id SERIAL PRIMARY KEY,
name VARCHAR(100),
ssn_encrypted TEXT,
credit_card_encrypted TEXT
);
-- Insert encrypted data
INSERT INTO secure_customers (name, ssn_encrypted, credit_card_encrypted)
VALUES (
'John Doe',
encrypt_sensitive_data('123-45-6789', 'mysecretkey'),
encrypt_sensitive_data('4111111111111111', 'mysecretkey')
);
35. How do you handle database sharding?
Answer: Database sharding distributes data across multiple databases to improve performance and scalability.
Horizontal Sharding Implementation:
-- Shard routing table
CREATE TABLE shard_routing (
shard_id INT PRIMARY KEY,
shard_name VARCHAR(50),
connection_string TEXT,
key_range_start BIGINT,
key_range_end BIGINT,
is_active BOOLEAN DEFAULT TRUE
);
-- Shard-aware user table
CREATE TABLE users (
user_id BIGINT PRIMARY KEY,
username VARCHAR(50),
email VARCHAR(100),
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- Shard routing function
CREATE OR REPLACE FUNCTION get_shard_for_user(user_id BIGINT)
RETURNS INT AS $$
DECLARE
shard_id INT;
BEGIN
SELECT sr.shard_id INTO shard_id
FROM shard_routing sr
WHERE user_id BETWEEN sr.key_range_start AND sr.key_range_end
AND sr.is_active = TRUE;
IF shard_id IS NULL THEN
RAISE EXCEPTION 'No shard found for user_id: %', user_id;
END IF;
RETURN shard_id;
END;
$$ LANGUAGE plpgsql;
-- Shard-aware insert procedure
CREATE OR REPLACE PROCEDURE insert_user_sharded(
p_user_id BIGINT,
p_username VARCHAR(50),
p_email VARCHAR(100)
)
LANGUAGE plpgsql
AS $$
DECLARE
target_shard INT;
shard_connection TEXT;
BEGIN
-- Determine target shard
SELECT get_shard_for_user(p_user_id) INTO target_shard;
-- Get connection string for shard
SELECT connection_string INTO shard_connection
FROM shard_routing
WHERE shard_id = target_shard;
-- Execute on target shard (simplified - would use dynamic SQL)
-- dblink or similar would be used in practice
RAISE NOTICE 'Inserting user % into shard %', p_user_id, target_shard;
END;
$$;
36. How do you implement database mirroring?
Answer: Database mirroring provides high availability by maintaining a synchronized copy of the database.
SQL Server Mirroring Setup:
-- Principal server setup
ALTER DATABASE MyDatabase
SET PARTNER = 'TCP://mirror_server:5022';
-- Mirror server setup
ALTER DATABASE MyDatabase
SET PARTNER = 'TCP://principal_server:5022';
-- Monitor mirroring status
SELECT
database_name,
mirroring_state_desc,
mirroring_role_desc,
mirroring_safety_level_desc
FROM sys.database_mirroring
WHERE database_id = DB_ID('MyDatabase');
-- Failover procedure
CREATE PROCEDURE PerformFailover
AS
BEGIN
-- Check if mirror is synchronized
IF EXISTS (
SELECT 1 FROM sys.database_mirroring
WHERE database_id = DB_ID('MyDatabase')
AND mirroring_state = 4 -- Synchronized
)
BEGIN
ALTER DATABASE MyDatabase SET PARTNER FAILOVER;
PRINT 'Failover completed successfully';
END
ELSE
BEGIN
RAISERROR('Mirror is not synchronized', 16, 1);
END
END;
37. How do you handle database failover scenarios?
Answer: Database failover ensures continuous availability when primary systems fail.
Always On Availability Groups (SQL Server):
-- Create availability group
CREATE AVAILABILITY GROUP [MyAG]
WITH (
AUTOMATED_BACKUP_PREFERENCE = PRIMARY,
FAILOVER_MODE = AUTOMATIC,
HEALTH_CHECK_TIMEOUT = 30000
)
FOR DATABASE [MyDatabase]
REPLICA ON
'primary_server' WITH (
ENDPOINT_URL = 'TCP://primary_server:5022',
AVAILABILITY_MODE = SYNCHRONOUS_COMMIT,
FAILOVER_MODE = AUTOMATIC,
SEEDING_MODE = AUTOMATIC
),
'secondary_server' WITH (
ENDPOINT_URL = 'TCP://secondary_server:5022',
AVAILABILITY_MODE = SYNCHRONOUS_COMMIT,
FAILOVER_MODE = AUTOMATIC,
SEEDING_MODE = AUTOMATIC
);
-- Monitor availability group
SELECT
ag.name AS availability_group,
ar.replica_server_name,
ars.role_desc,
ars.operational_state_desc,
ars.connected_state_desc
FROM sys.availability_groups ag
JOIN sys.availability_replicas ar ON ag.group_id = ar.group_id
JOIN sys.dm_hadr_availability_replica_states ars ON ar.replica_id = ars.replica_id;
PostgreSQL Failover with Patroni:
# patroni.yml configuration
scope: postgres-cluster
namespace: /db/
name: postgres-node-1
restapi:
listen: 0.0.0.0:8008
connect_address: 192.168.1.10:8008
etcd:
host: 192.168.1.5:2379
bootstrap:
dcs:
ttl: 30
loop_wait: 10
retry_timeout: 10
maximum_lag_on_failover: 1048576
postgresql:
use_pg_rewind: true
parameters:
max_connections: 100
shared_buffers: 256MB
wal_level: replica
hot_standby: "on"
max_wal_senders: 10
max_replication_slots: 10
wal_keep_segments: 8
initdb:
- encoding: UTF8
- data-checksums
postgresql:
listen: 0.0.0.0:5432
connect_address: 192.168.1.10:5432
data_dir: /var/lib/postgresql/9.6/main
pgpass: /tmp/pgpass
authentication:
replication:
username: replicator
password: rep-pass
superuser:
username: postgres
password: admin-pass
parameters:
unix_socket_directories: '.'
tags:
nofailover: false
noloadbalance: false
clonefrom: false
nosync: false
38. How do you implement data archival and purging strategies?
Answer: Data archival and purging strategies manage data lifecycle and compliance requirements.
Automated Archival Strategy:
-- Archive configuration table
CREATE TABLE archive_config (
table_name VARCHAR(128) PRIMARY KEY,
retention_days INT,
archive_table_name VARCHAR(128),
purge_after_days INT,
is_active BOOLEAN DEFAULT TRUE,
last_archive_date TIMESTAMP
);
-- Archive procedure
CREATE OR REPLACE PROCEDURE archive_old_data()
LANGUAGE plpgsql
AS $$
DECLARE
config RECORD;
archive_date DATE;
purge_date DATE;
BEGIN
FOR config IN
SELECT * FROM archive_config WHERE is_active = TRUE
LOOP
archive_date := CURRENT_DATE - config.retention_days;
purge_date := CURRENT_DATE - config.purge_after_days;
-- Archive old data
EXECUTE format(
'INSERT INTO %I SELECT * FROM %I WHERE created_date < $1',
config.archive_table_name,
config.table_name
) USING archive_date;
-- Delete archived data from main table
EXECUTE format(
'DELETE FROM %I WHERE created_date < $1',
config.table_name
) USING archive_date;
-- Purge very old archived data
EXECUTE format(
'DELETE FROM %I WHERE created_date < $1',
config.archive_table_name
) USING purge_date;
-- Update last archive date
UPDATE archive_config
SET last_archive_date = CURRENT_TIMESTAMP
WHERE table_name = config.table_name;
RAISE NOTICE 'Archived data from table: %', config.table_name;
END LOOP;
END;
$$;
-- Schedule archival job
SELECT cron.schedule(
'archive-old-data',
'0 2 * * *', -- Daily at 2 AM
'CALL archive_old_data()'
);
39. How do you handle database maintenance windows?
Answer: Database maintenance windows require careful planning to minimize downtime and impact.
Maintenance Window Management:
-- Maintenance window configuration
CREATE TABLE maintenance_windows (
window_id SERIAL PRIMARY KEY,
window_name VARCHAR(100),
start_time TIME,
end_time TIME,
days_of_week INTEGER[], -- 0=Sunday, 1=Monday, etc.
is_active BOOLEAN DEFAULT TRUE,
max_duration_minutes INTEGER DEFAULT 60
);
-- Maintenance tasks
CREATE TABLE maintenance_tasks (
task_id SERIAL PRIMARY KEY,
task_name VARCHAR(100),
task_type VARCHAR(50), -- 'VACUUM', 'ANALYZE', 'REINDEX', 'BACKUP'
target_table VARCHAR(128),
priority INTEGER DEFAULT 5,
estimated_duration_minutes INTEGER,
last_run TIMESTAMP,
next_run TIMESTAMP,
is_enabled BOOLEAN DEFAULT TRUE
);
-- Maintenance execution procedure
CREATE OR REPLACE PROCEDURE execute_maintenance_window()
LANGUAGE plpgsql
AS $$
DECLARE
current_window RECORD;
task RECORD;
start_time TIMESTAMP;
end_time TIMESTAMP;
BEGIN
-- Check if we're in a maintenance window
SELECT * INTO current_window
FROM maintenance_windows
WHERE is_active = TRUE
AND EXTRACT(DOW FROM CURRENT_TIMESTAMP) = ANY(days_of_week)
AND CURRENT_TIME BETWEEN start_time AND end_time;
IF current_window IS NULL THEN
RAISE NOTICE 'Not in maintenance window';
RETURN;
END IF;
start_time := CURRENT_TIMESTAMP;
end_time := start_time + (current_window.max_duration_minutes || ' minutes')::INTERVAL;
-- Execute maintenance tasks
FOR task IN
SELECT * FROM maintenance_tasks
WHERE is_enabled = TRUE
AND next_run <= CURRENT_TIMESTAMP
ORDER BY priority DESC, estimated_duration_minutes ASC
LOOP
-- Check if we have time remaining
IF CURRENT_TIMESTAMP + (task.estimated_duration_minutes || ' minutes')::INTERVAL > end_time THEN
RAISE NOTICE 'Maintenance window ending, stopping tasks';
EXIT;
END IF;
-- Execute task based on type
CASE task.task_type
WHEN 'VACUUM' THEN
EXECUTE format('VACUUM ANALYZE %I', task.target_table);
WHEN 'REINDEX' THEN
EXECUTE format('REINDEX TABLE %I', task.target_table);
WHEN 'ANALYZE' THEN
EXECUTE format('ANALYZE %I', task.target_table);
ELSE
RAISE NOTICE 'Unknown task type: %', task.task_type;
END CASE;
-- Update task status
UPDATE maintenance_tasks
SET last_run = CURRENT_TIMESTAMP,
next_run = CURRENT_TIMESTAMP + '1 day'::INTERVAL
WHERE task_id = task.task_id;
RAISE NOTICE 'Completed task: %', task.task_name;
END LOOP;
END;
$$;
40. How do you implement database monitoring and alerting?
Answer: Database monitoring and alerting systems track performance, availability, and health metrics.
Comprehensive Monitoring System:
-- Monitoring configuration
CREATE TABLE monitoring_config (
metric_id SERIAL PRIMARY KEY,
metric_name VARCHAR(100),
metric_query TEXT,
threshold_value NUMERIC,
threshold_operator VARCHAR(10), -- '>', '<', '=', '!='
check_interval_minutes INTEGER DEFAULT 5,
is_active BOOLEAN DEFAULT TRUE,
alert_email VARCHAR(255),
severity VARCHAR(20) DEFAULT 'WARNING' -- 'INFO', 'WARNING', 'CRITICAL'
);
-- Monitoring results
CREATE TABLE monitoring_results (
result_id SERIAL PRIMARY KEY,
metric_id INTEGER REFERENCES monitoring_config(metric_id),
check_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
metric_value NUMERIC,
is_alert BOOLEAN,
alert_message TEXT
);
-- Monitoring procedure
CREATE OR REPLACE PROCEDURE run_database_monitoring()
LANGUAGE plpgsql
AS $$
DECLARE
metric RECORD;
result_value NUMERIC;
is_alert BOOLEAN := FALSE;
alert_msg TEXT;
BEGIN
FOR metric IN
SELECT * FROM monitoring_config
WHERE is_active = TRUE
AND (last_check IS NULL OR
last_check < CURRENT_TIMESTAMP - (check_interval_minutes || ' minutes')::INTERVAL)
LOOP
-- Execute metric query
BEGIN
EXECUTE metric.metric_query INTO result_value;
-- Check threshold
EXECUTE format(
'SELECT $1 %s $2',
metric.threshold_operator
) INTO is_alert USING result_value, metric.threshold_value;
-- Generate alert message
IF is_alert THEN
alert_msg := format(
'Metric %s exceeded threshold. Value: %s, Threshold: %s %s',
metric.metric_name,
result_value,
metric.threshold_operator,
metric.threshold_value
);
END IF;
-- Store result
INSERT INTO monitoring_results (
metric_id, metric_value, is_alert, alert_message
) VALUES (
metric.metric_id, result_value, is_alert, alert_msg
);
-- Send alert if needed
IF is_alert THEN
PERFORM send_alert_email(metric.alert_email, alert_msg, metric.severity);
END IF;
EXCEPTION WHEN OTHERS THEN
-- Log monitoring error
INSERT INTO monitoring_results (
metric_id, metric_value, is_alert, alert_message
) VALUES (
metric.metric_id, NULL, TRUE,
'Monitoring query failed: ' || SQLERRM
);
END;
-- Update last check time
UPDATE monitoring_config
SET last_check = CURRENT_TIMESTAMP
WHERE metric_id = metric.metric_id;
END LOOP;
END;
$$;
-- Example monitoring queries
INSERT INTO monitoring_config (metric_name, metric_query, threshold_value, threshold_operator) VALUES
('Active Connections', 'SELECT count(*) FROM pg_stat_activity WHERE state = ''active''', 100, '>'),
('Database Size (GB)', 'SELECT pg_database_size(current_database()) / 1024.0 / 1024.0 / 1024.0', 50, '>'),
('Slow Queries (>5s)', 'SELECT count(*) FROM pg_stat_activity WHERE state = ''active'' AND query_start < now() - interval ''5 seconds''', 10, '>'),
('Lock Count', 'SELECT count(*) FROM pg_locks WHERE NOT granted', 5, '>'),
('WAL Files', 'SELECT count(*) FROM pg_ls_waldir()', 100, '>');
Database Security (Questions 41-50)
41. SQL Injection Prevention
Answer: SQL injection prevention requires multiple layers of defense using parameterized queries, input validation, and proper escaping.
Implementation Examples:
-- BAD: Vulnerable to SQL injection
SELECT * FROM users WHERE username = '$username' AND password = '$password';
-- GOOD: Parameterized query
SELECT * FROM users WHERE username = ? AND password = ?;
Python with parameterized queries
import sqlite3
def secure_login(username, password): conn = sqlite3.connect('database.db') cursor = conn.cursor()
# Parameterized query prevents SQL injection
query = "SELECT * FROM users WHERE username = ? AND password = ?"
cursor.execute(query, (username, password))
result = cursor.fetchone()
conn.close()
return result
// C# with Entity Framework
public User GetUser(string username, string password)
{
// Entity Framework automatically uses parameterized queries
return context.Users
.Where(u => u.Username == username && u.Password == password)
.FirstOrDefault();
}
42. Database User Permissions and Roles
Answer: Implement principle of least privilege using role-based access control (RBAC) with granular permissions.
-- Create roles
CREATE ROLE read_only_role;
CREATE ROLE data_analyst_role;
CREATE ROLE admin_role;
-- Grant permissions to roles
GRANT SELECT ON ALL TABLES IN SCHEMA public TO read_only_role;
GRANT SELECT, INSERT, UPDATE ON customer_data TO data_analyst_role;
GRANT ALL PRIVILEGES ON ALL TABLES IN SCHEMA public TO admin_role;
-- Create users and assign roles
CREATE USER analyst1 WITH PASSWORD 'secure_password';
GRANT data_analyst_role TO analyst1;
-- Application-specific user
CREATE USER app_user WITH PASSWORD 'app_password';
GRANT CONNECT ON DATABASE mydb TO app_user;
GRANT USAGE ON SCHEMA public TO app_user;
GRANT SELECT, INSERT, UPDATE ON specific_tables TO app_user;
-- Row-level security (PostgreSQL)
CREATE POLICY user_data_policy ON customer_data
FOR ALL
USING (created_by = current_user);
ALTER TABLE customer_data ENABLE ROW LEVEL SECURITY;
43. Audit Logging for Sensitive Operations
Answer: Implement comprehensive audit trails using triggers, application logging, and dedicated audit tables.
-- Create audit table
CREATE TABLE audit_log (
id SERIAL PRIMARY KEY,
table_name VARCHAR(100),
operation VARCHAR(10),
old_data JSONB,
new_data JSONB,
user_id VARCHAR(100),
timestamp TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
ip_address INET
);
-- Audit trigger function
CREATE OR REPLACE FUNCTION audit_trigger_function()
RETURNS TRIGGER AS $$
BEGIN
IF TG_OP = 'INSERT' THEN
INSERT INTO audit_log (table_name, operation, new_data, user_id)
VALUES (TG_TABLE_NAME, 'INSERT', to_jsonb(NEW), current_user);
RETURN NEW;
ELSIF TG_OP = 'UPDATE' THEN
INSERT INTO audit_log (table_name, operation, old_data, new_data, user_id)
VALUES (TG_TABLE_NAME, 'UPDATE', to_jsonb(OLD), to_jsonb(NEW), current_user);
RETURN NEW;
ELSIF TG_OP = 'DELETE' THEN
INSERT INTO audit_log (table_name, operation, old_data, user_id)
VALUES (TG_TABLE_NAME, 'DELETE', to_jsonb(OLD), current_user);
RETURN OLD;
END IF;
RETURN NULL;
END;
$$ LANGUAGE plpgsql;
-- Apply trigger to sensitive tables
CREATE TRIGGER audit_customer_data_trigger
AFTER INSERT OR UPDATE OR DELETE ON customer_data
FOR EACH ROW EXECUTE FUNCTION audit_trigger_function();
Application-level audit logging
import logging from datetime import datetime
class AuditLogger: def init(self): self.logger = logging.getLogger('audit') self.logger.setLevel(logging.INFO)
def log_operation(self, user_id, operation, table, record_id, details):
audit_entry = {
'timestamp': datetime.utcnow().isoformat(),
'user_id': user_id,
'operation': operation,
'table': table,
'record_id': record_id,
'details': details,
'ip_address': self.get_client_ip()
}
self.logger.info(f"AUDIT: {audit_entry}")
44. Database Encryption
Answer: Implement encryption at rest, in transit, and for sensitive data fields.
-- PostgreSQL with pgcrypto for field-level encryption
CREATE EXTENSION pgcrypto;
-- Encrypt sensitive data
CREATE TABLE users (
id SERIAL PRIMARY KEY,
username VARCHAR(100),
email VARCHAR(255),
encrypted_ssn BYTEA,
encrypted_credit_card BYTEA
);
-- Function to encrypt data
CREATE OR REPLACE FUNCTION encrypt_sensitive_data(data TEXT, key TEXT)
RETURNS BYTEA AS $$
BEGIN
RETURN pgp_sym_encrypt(data, key);
END;
$$ LANGUAGE plpgsql;
-- Function to decrypt data
CREATE OR REPLACE FUNCTION decrypt_sensitive_data(encrypted_data BYTEA, key TEXT)
RETURNS TEXT AS $$
BEGIN
RETURN pgp_sym_decrypt(encrypted_data, key);
END;
$$ LANGUAGE plpgsql;
Application-level encryption
from cryptography.fernet import Fernet import base64
class DatabaseEncryption: def init(self, key): self.cipher = Fernet(key)
def encrypt_field(self, data):
if data is None:
return None
return self.cipher.encrypt(data.encode()).decode()
def decrypt_field(self, encrypted_data):
if encrypted_data is None:
return None
return self.cipher.decrypt(encrypted_data.encode()).decode()
45. Database Firewall Rules
Answer: Implement network-level and application-level firewall rules to control database access.
-- PostgreSQL pg_hba.conf configuration
TYPE DATABASE USER ADDRESS METHOD
local all all md5 host myapp_db app_user 192.168.1.0/24 md5 host myapp_db app_user 10.0.0.0/8 md5 host all all 0.0.0.0/0 reject
Application-level firewall
import ipaddress from functools import wraps
class DatabaseFirewall: def init(self): self.allowed_ips = [ ipaddress.ip_network('192.168.1.0/24'), ipaddress.ip_network('10.0.0.0/8') ]
def check_ip_access(self, client_ip):
client_ip_obj = ipaddress.ip_address(client_ip)
return any(client_ip_obj in network for network in self.allowed_ips)
def firewall_decorator(self, func):
@wraps(func)
def wrapper(*args, **kwargs):
client_ip = self.get_client_ip()
if not self.check_ip_access(client_ip):
raise SecurityException(f"Access denied from IP: {client_ip}")
return func(*args, **kwargs)
return wrapper
46. Database Access from Applications
Answer: Use connection pooling, connection string security, and proper authentication.
Secure database connection with connection pooling
import psycopg2 from psycopg2 import pool from contextlib import contextmanager
class DatabaseConnectionManager: def init(self, config): self.pool = psycopg2.pool.ThreadedConnectionPool( minconn=5, maxconn=20, host=config['host'], database=config['database'], user=config['user'], password=config['password'], sslmode='require' )
@contextmanager
def get_connection(self):
conn = self.pool.getconn()
try:
yield conn
finally:
self.pool.putconn(conn)
def execute_query(self, query, params=None):
with self.get_connection() as conn:
with conn.cursor() as cursor:
cursor.execute(query, params)
return cursor.fetchall()
// C# with connection string security
public class DatabaseService
{
private readonly string _connectionString;
public DatabaseService(IConfiguration config)
{
_connectionString = config.GetConnectionString("DefaultConnection");
}
public async Task<T> ExecuteQueryAsync<T>(string query, object parameters = null)
{
using var connection = new SqlConnection(_connectionString);
await connection.OpenAsync();
// Use Dapper for safe parameterized queries
return await connection.QueryFirstOrDefaultAsync<T>(query, parameters);
}
}
47. Data Masking for Sensitive Information
Answer: Implement dynamic data masking and static data masking for sensitive data protection.
-- Dynamic data masking (SQL Server)
CREATE TABLE customer_data (
id INT PRIMARY KEY,
name VARCHAR(100),
email VARCHAR(255) MASKED WITH (FUNCTION = 'email()'),
ssn VARCHAR(11) MASKED WITH (FUNCTION = 'partial(0, "XXX-XX-", 4)'),
credit_card VARCHAR(16) MASKED WITH (FUNCTION = 'partial(0, "XXXX-XXXX-XXXX-", 4)')
);
-- Grant unmask permission to authorized users
GRANT UNMASK ON customer_data TO authorized_user;
-- PostgreSQL data masking function
CREATE OR REPLACE FUNCTION mask_email(email TEXT)
RETURNS TEXT AS $$
BEGIN
IF email IS NULL THEN
RETURN NULL;
END IF;
RETURN substring(email from 1 for 2) || '***@' ||
substring(email from position('@' in email) + 1);
END;
$$ LANGUAGE plpgsql;
-- Create view with masked data
CREATE VIEW masked_customer_view AS
SELECT
id,
name,
mask_email(email) as email,
'***-**-' || substring(ssn from 8) as ssn
FROM customer_data;
48. Database Security Compliance
Answer: Implement compliance frameworks like GDPR, HIPAA, SOX, and PCI DSS.
-- GDPR compliance: Data retention and right to be forgotten
CREATE TABLE data_retention_policy (
table_name VARCHAR(100),
retention_period_days INT,
deletion_strategy VARCHAR(50)
);
-- Function to implement right to be forgotten
CREATE OR REPLACE FUNCTION delete_user_data(user_id INT)
RETURNS VOID AS $$
BEGIN
-- Anonymize instead of delete for audit purposes
UPDATE customer_data
SET name = 'DELETED',
email = 'deleted@example.com',
ssn = NULL,
deleted_at = CURRENT_TIMESTAMP
WHERE id = user_id;
-- Log the deletion
INSERT INTO audit_log (table_name, operation, user_id, details)
VALUES ('customer_data', 'GDPR_DELETE', user_id, 'Right to be forgotten');
END;
$$ LANGUAGE plpgsql;
Compliance monitoring
class ComplianceMonitor: def init(self): self.compliance_rules = { 'gdpr': self.check_gdpr_compliance, 'hipaa': self.check_hipaa_compliance, 'pci': self.check_pci_compliance }
def check_gdpr_compliance(self, data):
# Check for data minimization
# Verify consent tracking
# Ensure right to be forgotten
pass
def generate_compliance_report(self):
report = {
'timestamp': datetime.utcnow(),
'gdpr_status': self.check_gdpr_compliance(),
'hipaa_status': self.check_hipaa_compliance(),
'pci_status': self.check_pci_compliance()
}
return report
49. Database Vulnerability Scanning
Answer: Implement automated vulnerability scanning and security assessments.
Database vulnerability scanner
import subprocess import json from typing import List, Dict
class DatabaseVulnerabilityScanner: def init(self, db_config): self.db_config = db_config
def scan_for_vulnerabilities(self) -> List[Dict]:
vulnerabilities = []
# Check for weak passwords
vulnerabilities.extend(self.check_weak_passwords())
# Check for excessive privileges
vulnerabilities.extend(self.check_excessive_privileges())
# Check for unpatched vulnerabilities
vulnerabilities.extend(self.check_database_version())
# Check for misconfigurations
vulnerabilities.extend(self.check_misconfigurations())
return vulnerabilities
def check_weak_passwords(self):
# Implementation for password strength checking
pass
def check_excessive_privileges(self):
query = """
SELECT grantee, privilege_type, table_name
FROM information_schema.role_table_grants
WHERE privilege_type IN ('ALL', 'DELETE', 'DROP')
"""
# Execute and analyze results
pass
50. Database Security Incident Response
Answer: Implement incident response procedures with automated detection and response.
Security incident response system
import logging from datetime import datetime from typing import Dict, List
class SecurityIncidentResponse: def init(self): self.logger = logging.getLogger('security_incident') self.incident_handlers = { 'sql_injection': self.handle_sql_injection, 'unauthorized_access': self.handle_unauthorized_access, 'data_breach': self.handle_data_breach }
def detect_incident(self, event_type: str, details: Dict):
if self.is_security_incident(event_type, details):
self.respond_to_incident(event_type, details)
def respond_to_incident(self, incident_type: str, details: Dict):
# Log the incident
self.log_incident(incident_type, details)
# Execute response plan
if incident_type in self.incident_handlers:
self.incident_handlers[incident_type](details)
# Notify stakeholders
self.notify_stakeholders(incident_type, details)
def handle_sql_injection(self, details: Dict):
# Block suspicious IP
self.block_ip(details['source_ip'])
# Quarantine affected data
self.quarantine_data(details['affected_tables'])
# Generate incident report
self.generate_incident_report('sql_injection', details)
Advanced Database Features (Questions 51-60)
51. Database Triggers Effectively
Answer: Use triggers for data integrity, audit trails, and business logic enforcement.
-- Business logic trigger
CREATE OR REPLACE FUNCTION update_order_status()
RETURNS TRIGGER AS $$
BEGIN
-- Update order status based on payment status
IF NEW.payment_status = 'PAID' THEN
NEW.order_status = 'CONFIRMED';
NEW.confirmed_at = CURRENT_TIMESTAMP;
ELSIF NEW.payment_status = 'FAILED' THEN
NEW.order_status = 'CANCELLED';
NEW.cancelled_at = CURRENT_TIMESTAMP;
END IF;
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER order_status_trigger
BEFORE UPDATE ON orders
FOR EACH ROW EXECUTE FUNCTION update_order_status();
-- Data validation trigger
CREATE OR REPLACE FUNCTION validate_customer_data()
RETURNS TRIGGER AS $$
BEGIN
-- Validate email format
IF NEW.email !~ '^[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Za-z]{2,}$' THEN
RAISE EXCEPTION 'Invalid email format';
END IF;
-- Validate age
IF NEW.age < 0 OR NEW.age > 150 THEN
RAISE EXCEPTION 'Invalid age';
END IF;
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
52. Window Functions for Advanced Analytics
Answer: Use window functions for ranking, running totals, and complex analytical queries.
-- Sales analysis with window functions
SELECT
product_name,
sales_date,
daily_sales,
-- Running total
SUM(daily_sales) OVER (
PARTITION BY product_name
ORDER BY sales_date
ROWS UNBOUNDED PRECEDING
) as running_total,
-- Moving average
AVG(daily_sales) OVER (
PARTITION BY product_name
ORDER BY sales_date
ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
) as moving_avg_7d,
-- Rank within category
RANK() OVER (
PARTITION BY category
ORDER BY daily_sales DESC
) as category_rank,
-- Percentile
PERCENT_RANK() OVER (
PARTITION BY category
ORDER BY daily_sales
) as sales_percentile
FROM sales_data;
-- Customer segmentation
SELECT
customer_id,
total_purchases,
NTILE(4) OVER (ORDER BY total_purchases DESC) as customer_segment,
LAG(total_purchases) OVER (
PARTITION BY customer_id
ORDER BY purchase_date
) as previous_purchase,
LEAD(total_purchases) OVER (
PARTITION BY customer_id
ORDER BY purchase_date
) as next_purchase
FROM customer_purchases;
53. Database Views for Security and Abstraction
Answer: Use views to provide controlled access and simplify complex queries.
-- Security view for customer data
CREATE VIEW customer_public_view AS
SELECT
id,
first_name,
last_name,
email,
phone,
created_date
FROM customers
WHERE active = true;
-- Grant access to view instead of table
GRANT SELECT ON customer_public_view TO public_role;
REVOKE ALL ON customers FROM public_role;
-- Complex business logic view
CREATE VIEW sales_summary_view AS
SELECT
c.customer_id,
c.customer_name,
COUNT(o.order_id) as total_orders,
SUM(o.order_amount) as total_spent,
AVG(o.order_amount) as avg_order_value,
MAX(o.order_date) as last_order_date,
CASE
WHEN SUM(o.order_amount) > 10000 THEN 'VIP'
WHEN SUM(o.order_amount) > 5000 THEN 'Premium'
ELSE 'Standard'
END as customer_tier
FROM customers c
LEFT JOIN orders o ON c.customer_id = o.customer_id
WHERE o.order_date >= CURRENT_DATE - INTERVAL '1 year'
GROUP BY c.customer_id, c.customer_name;
54. Database Transactions and Isolation Levels
Answer: Choose appropriate isolation levels based on consistency and performance requirements.
-- Transaction with proper isolation level
BEGIN TRANSACTION ISOLATION LEVEL SERIALIZABLE;
-- Critical financial transaction
UPDATE accounts
SET balance = balance - 1000
WHERE account_id = 'A001' AND balance >= 1000;
IF ROW_COUNT = 0 THEN
ROLLBACK;
RAISE EXCEPTION 'Insufficient funds';
END IF;
UPDATE accounts
SET balance = balance + 1000
WHERE account_id = 'A002';
COMMIT;
-- Read-committed for reporting
BEGIN TRANSACTION ISOLATION LEVEL READ COMMITTED;
SELECT
account_id,
balance,
last_transaction_date
FROM accounts
WHERE account_type = 'SAVINGS';
COMMIT;
Application-level transaction management
import psycopg2 from contextlib import contextmanager
class TransactionManager: def init(self, connection): self.connection = connection
@contextmanager
def transaction(self, isolation_level=None):
if isolation_level:
self.connection.set_isolation_level(isolation_level)
try:
yield self.connection
self.connection.commit()
except Exception:
self.connection.rollback()
raise
def execute_with_retry(self, query, params=None, max_retries=3):
for attempt in range(max_retries):
try:
with self.transaction() as conn:
with conn.cursor() as cursor:
cursor.execute(query, params)
return cursor.fetchall()
except psycopg2.extensions.TransactionRollbackError:
if attempt == max_retries - 1:
raise
continue
55. Database Stored Procedures with Error Handling
Answer: Implement robust stored procedures with comprehensive error handling.
-- Stored procedure with error handling
CREATE OR REPLACE FUNCTION process_order(
p_customer_id INT,
p_product_id INT,
p_quantity INT
) RETURNS JSON AS $$
DECLARE
v_order_id INT;
v_product_price DECIMAL(10,2);
v_total_amount DECIMAL(10,2);
v_customer_balance DECIMAL(10,2);
v_result JSON;
BEGIN
-- Start transaction
BEGIN
-- Validate customer exists
IF NOT EXISTS (SELECT 1 FROM customers WHERE customer_id = p_customer_id) THEN
RAISE EXCEPTION 'Customer not found: %', p_customer_id;
END IF;
-- Get product price
SELECT price INTO v_product_price
FROM products
WHERE product_id = p_product_id;
IF NOT FOUND THEN
RAISE EXCEPTION 'Product not found: %', p_product_id;
END IF;
-- Calculate total
v_total_amount := v_product_price * p_quantity;
-- Check customer balance
SELECT balance INTO v_customer_balance
FROM customers
WHERE customer_id = p_customer_id;
IF v_customer_balance < v_total_amount THEN
RAISE EXCEPTION 'Insufficient balance. Required: %, Available: %',
v_total_amount, v_customer_balance;
END IF;
-- Create order
INSERT INTO orders (customer_id, product_id, quantity, total_amount, order_date)
VALUES (p_customer_id, p_product_id, p_quantity, v_total_amount, CURRENT_TIMESTAMP)
RETURNING order_id INTO v_order_id;
-- Update customer balance
UPDATE customers
SET balance = balance - v_total_amount
WHERE customer_id = p_customer_id;
-- Update product inventory
UPDATE products
SET stock_quantity = stock_quantity - p_quantity
WHERE product_id = p_product_id;
-- Return success result
v_result := json_build_object(
'success', true,
'order_id', v_order_id,
'total_amount', v_total_amount,
'message', 'Order processed successfully'
);
RETURN v_result;
EXCEPTION
WHEN OTHERS THEN
-- Log error
INSERT INTO error_log (error_message, error_detail, created_at)
VALUES (SQLERRM, SQLSTATE, CURRENT_TIMESTAMP);
-- Return error result
v_result := json_build_object(
'success', false,
'error', SQLERRM,
'error_code', SQLSTATE
);
RETURN v_result;
END;
END;
$$ LANGUAGE plpgsql;
56. Database Functions for Code Reuse
Answer: Create reusable functions for common operations and business logic.
-- Utility functions
CREATE OR REPLACE FUNCTION calculate_age(birth_date DATE)
RETURNS INT AS $$
BEGIN
RETURN EXTRACT(YEAR FROM AGE(birth_date));
END;
$$ LANGUAGE plpgsql;
-- Business logic functions
CREATE OR REPLACE FUNCTION calculate_discount(
p_customer_tier VARCHAR(20),
p_order_amount DECIMAL(10,2)
) RETURNS DECIMAL(10,2) AS $$
BEGIN
RETURN CASE p_customer_tier
WHEN 'VIP' THEN p_order_amount * 0.15
WHEN 'Premium' THEN p_order_amount * 0.10
WHEN 'Standard' THEN p_order_amount * 0.05
ELSE 0
END;
END;
$$ LANGUAGE plpgsql;
-- Data validation functions
CREATE OR REPLACE FUNCTION is_valid_email(email TEXT)
RETURNS BOOLEAN AS $$
BEGIN
RETURN email ~ '^[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Za-z]{2,}$';
END;
$$ LANGUAGE plpgsql;
-- Complex business function
CREATE OR REPLACE FUNCTION get_customer_recommendations(p_customer_id INT)
RETURNS TABLE (
product_id INT,
product_name VARCHAR(100),
recommendation_score DECIMAL(5,4)
) AS $$
BEGIN
RETURN QUERY
SELECT
p.product_id,
p.product_name,
(
-- Purchase history similarity
COUNT(DISTINCT o2.customer_id) * 0.4 +
-- Category preference
CASE WHEN p.category_id = (
SELECT category_id
FROM products
WHERE product_id = (
SELECT product_id
FROM order_items oi
JOIN orders o ON oi.order_id = o.order_id
WHERE o.customer_id = p_customer_id
GROUP BY oi.product_id
ORDER BY COUNT(*) DESC
LIMIT 1
)
) THEN 0.3 ELSE 0 END +
-- Price range preference
CASE WHEN p.price BETWEEN
(SELECT AVG(price) - STDDEV(price) FROM products) AND
(SELECT AVG(price) + STDDEV(price) FROM products)
THEN 0.3 ELSE 0 END
) as recommendation_score
FROM products p
LEFT JOIN order_items oi ON p.product_id = oi.product_id
LEFT JOIN orders o2 ON oi.order_id = o2.order_id
WHERE o2.customer_id != p_customer_id
GROUP BY p.product_id, p.product_name
ORDER BY recommendation_score DESC
LIMIT 10;
END;
$$ LANGUAGE plpgsql;
57. Database Cursors vs Set-Based Operations
Answer: Use set-based operations for performance, cursors only when necessary for complex logic.
-- BAD: Cursor-based approach
CREATE OR REPLACE FUNCTION update_customer_tiers_cursor()
RETURNS VOID AS $$
DECLARE
customer_record RECORD;
customer_cursor CURSOR FOR
SELECT customer_id, total_spent
FROM customer_summary;
BEGIN
OPEN customer_cursor;
LOOP
FETCH customer_cursor INTO customer_record;
EXIT WHEN NOT FOUND;
UPDATE customers
SET tier = CASE
WHEN customer_record.total_spent > 10000 THEN 'VIP'
WHEN customer_record.total_spent > 5000 THEN 'Premium'
ELSE 'Standard'
END
WHERE customer_id = customer_record.customer_id;
END LOOP;
CLOSE customer_cursor;
END;
$$ LANGUAGE plpgsql;
-- GOOD: Set-based approach
CREATE OR REPLACE FUNCTION update_customer_tiers_set()
RETURNS VOID AS $$
BEGIN
UPDATE customers
SET tier = CASE
WHEN total_spent > 10000 THEN 'VIP'
WHEN total_spent > 5000 THEN 'Premium'
ELSE 'Standard'
END
FROM customer_summary cs
WHERE customers.customer_id = cs.customer_id;
END;
$$ LANGUAGE plpgsql;
-- When cursors are necessary (complex business logic)
CREATE OR REPLACE FUNCTION process_complex_business_logic()
RETURNS VOID AS $$
DECLARE
order_record RECORD;
order_cursor CURSOR FOR
SELECT o.order_id, o.customer_id, o.total_amount,
c.credit_score, c.payment_history
FROM orders o
JOIN customers c ON o.customer_id = c.customer_id
WHERE o.status = 'PENDING_APPROVAL';
BEGIN
OPEN order_cursor;
LOOP
FETCH order_cursor INTO order_record;
EXIT WHEN NOT FOUND;
-- Complex business logic that can't be done in a single query
IF order_record.credit_score < 600 AND order_record.total_amount > 1000 THEN
-- Require additional verification
UPDATE orders
SET status = 'REQUIRES_VERIFICATION',
verification_required = true
WHERE order_id = order_record.order_id;
ELSIF order_record.payment_history = 'GOOD' AND order_record.total_amount < 500 THEN
-- Auto-approve
UPDATE orders
SET status = 'APPROVED',
approved_at = CURRENT_TIMESTAMP
WHERE order_id = order_record.order_id;
ELSE
-- Manual review required
UPDATE orders
SET status = 'MANUAL_REVIEW'
WHERE order_id = order_record.order_id;
END IF;
END LOOP;
CLOSE order_cursor;
END;
$$ LANGUAGE plpgsql;
58. Temporary Tables vs Table Variables
Answer: Choose based on data size, persistence needs, and performance requirements.
-- Temporary tables for large datasets and complex operations
CREATE TEMPORARY TABLE temp_sales_analysis (
product_id INT,
product_name VARCHAR(100),
total_sales DECIMAL(10,2),
avg_price DECIMAL(10,2),
sales_count INT
) ON COMMIT DROP;
-- Populate temporary table
INSERT INTO temp_sales_analysis
SELECT
p.product_id,
p.product_name,
SUM(oi.quantity * oi.unit_price) as total_sales,
AVG(oi.unit_price) as avg_price,
COUNT(*) as sales_count
FROM products p
JOIN order_items oi ON p.product_id = oi.product_id
JOIN orders o ON oi.order_id = o.order_id
WHERE o.order_date >= CURRENT_DATE - INTERVAL '30 days'
GROUP BY p.product_id, p.product_name;
-- Use temporary table for complex analysis
SELECT
tsa.*,
RANK() OVER (ORDER BY tsa.total_sales DESC) as sales_rank,
PERCENT_RANK() OVER (ORDER BY tsa.total_sales) as sales_percentile
FROM temp_sales_analysis tsa
WHERE tsa.total_sales > (SELECT AVG(total_sales) FROM temp_sales_analysis);
-- Table variables for small datasets (SQL Server example)
DECLARE @CustomerIDs TABLE (
customer_id INT PRIMARY KEY
);
INSERT INTO @CustomerIDs (customer_id)
SELECT customer_id
FROM customers
WHERE registration_date >= '2024-01-01';
-- Use table variable for filtering
SELECT c.*, o.total_orders
FROM customers c
JOIN @CustomerIDs cid ON c.customer_id = cid.customer_id
LEFT JOIN (
SELECT customer_id, COUNT(*) as total_orders
FROM orders
GROUP BY customer_id
) o ON c.customer_id = o.customer_id;
59. Cross-Database Queries
Answer: Use linked servers, database links, or federation for cross-database operations.
-- PostgreSQL foreign data wrapper
CREATE EXTENSION postgres_fdw;
-- Create foreign server
CREATE SERVER remote_db_server
FOREIGN DATA WRAPPER postgres_fdw
OPTIONS (host 'remote-server.com', port '5432', dbname 'remote_db');
-- Create user mapping
CREATE USER MAPPING FOR current_user
SERVER remote_db_server
OPTIONS (user 'remote_user', password 'remote_password');
-- Create foreign table
CREATE FOREIGN TABLE remote_customers (
customer_id INT,
customer_name VARCHAR(100),
email VARCHAR(255)
) SERVER remote_db_server
OPTIONS (schema_name 'public', table_name 'customers');
-- Cross-database query
SELECT
lc.customer_id,
lc.customer_name,
rc.email,
lc.local_data
FROM local_customers lc
JOIN remote_customers rc ON lc.customer_id = rc.customer_id
WHERE lc.active = true;
-- SQL Server linked server
-- Create linked server
EXEC sp_addlinkedserver
@server = 'RemoteServer',
@srvproduct = '',
@provider = 'SQLNCLI',
@datasrc = 'remote-server.com';
-- Create linked server login
EXEC sp_addlinkedsrvlogin
@rmtsrvname = 'RemoteServer',
@useself = 'FALSE',
@rmtuser = 'remote_user',
@rmtpassword = 'remote_password';
-- Cross-database query
SELECT
lc.customer_id,
lc.customer_name,
rc.email,
lc.local_data
FROM local_database.dbo.local_customers lc
JOIN RemoteServer.remote_database.dbo.remote_customers rc
ON lc.customer_id = rc.customer_id
WHERE lc.active = 1;
60. Database Linked Servers
Answer: Configure and manage linked servers for distributed database operations.
-- SQL Server linked server configuration
-- Create linked server for Oracle
EXEC sp_addlinkedserver
@server = 'OracleServer',
@srvproduct = 'Oracle',
@provider = 'OraOLEDB.Oracle',
@datasrc = 'ORCL';
-- Create linked server for MySQL
EXEC sp_addlinkedserver
@server = 'MySQLServer',
@srvproduct = 'MySQL',
@provider = 'MSDASQL',
@datasrc = 'MySQL_DSN';
-- Create linked server login
EXEC sp_addlinkedsrvlogin
@rmtsrvname = 'OracleServer',
@useself = 'FALSE',
@rmtuser = 'oracle_user',
@rmtpassword = 'oracle_password';
-- Distributed query with error handling
BEGIN TRY
-- Query across multiple databases
SELECT
'SQL_Server' as source,
customer_id,
customer_name,
email
FROM local_customers
WHERE active = 1
UNION ALL
SELECT
'Oracle' as source,
customer_id,
customer_name,
email
FROM OracleServer..SCHEMA.CUSTOMERS
WHERE status = 'ACTIVE'
UNION ALL
SELECT
'MySQL' as source,
customer_id,
customer_name,
email
FROM MySQLServer..database.customers
WHERE is_active = 1;
END TRY
BEGIN CATCH
-- Log error and continue with available data
INSERT INTO error_log (error_message, error_detail, created_at)
VALUES (ERROR_MESSAGE(), ERROR_LINE(), GETDATE());
-- Return only local data if remote queries fail
SELECT
'SQL_Server' as source,
customer_id,
customer_name,
email
FROM local_customers
WHERE active = 1;
END CATCH
-- PostgreSQL foreign data wrapper for multiple databases
-- Oracle foreign data wrapper
CREATE EXTENSION oracle_fdw;
CREATE SERVER oracle_server
FOREIGN DATA WRAPPER oracle_fdw
OPTIONS (dbserver '//oracle-server:1521/ORCL');
-- MySQL foreign data wrapper
CREATE EXTENSION mysql_fdw;
CREATE SERVER mysql_server
FOREIGN DATA WRAPPER mysql_fdw
OPTIONS (host 'mysql-server.com', port '3306');
-- Create foreign tables
CREATE FOREIGN TABLE oracle_customers (
customer_id NUMBER,
customer_name VARCHAR2(100),
email VARCHAR2(255)
) SERVER oracle_server
OPTIONS (schema 'SCHEMA', table 'CUSTOMERS');
CREATE FOREIGN TABLE mysql_customers (
customer_id INT,
customer_name VARCHAR(100),
email VARCHAR(255)
) SERVER mysql_server
OPTIONS (database 'database', table 'customers');
-- Distributed query with materialized view for performance
CREATE MATERIALIZED VIEW distributed_customer_view AS
SELECT
'postgresql' as source,
customer_id,
customer_name,
email
FROM local_customers
WHERE active = true
UNION ALL
SELECT
'oracle' as source,
customer_id,
customer_name,
email
FROM oracle_customers
WHERE status = 'ACTIVE'
UNION ALL
SELECT
'mysql' as source,
customer_id,
customer_name,
email
FROM mysql_customers
WHERE is_active = 1;
-- Refresh materialized view periodically
REFRESH MATERIALIZED VIEW distributed_customer_view;
Database Performance & Operations (Questions 61-70)
61. How do you monitor database performance metrics?
Answer: Database performance monitoring involves tracking key metrics across multiple dimensions:
Key Metrics to Monitor: - Query Performance: Execution time, CPU usage, I/O operations - Resource Utilization: CPU, memory, disk I/O, network - Connection Metrics: Active connections, connection pool usage - Lock Contention: Lock wait times, deadlock frequency - Storage Metrics: Disk space, I/O latency, throughput
Implementation Example:
-- Create a performance monitoring table
CREATE TABLE performance_metrics (
id BIGINT IDENTITY(1,1) PRIMARY KEY,
metric_name VARCHAR(100),
metric_value DECIMAL(18,4),
metric_unit VARCHAR(20),
collection_time DATETIME2 DEFAULT GETDATE(),
server_name VARCHAR(100),
database_name VARCHAR(100)
);
-- Stored procedure to collect performance metrics
CREATE PROCEDURE CollectPerformanceMetrics
AS
BEGIN
-- CPU Usage
INSERT INTO performance_metrics (metric_name, metric_value, metric_unit, server_name, database_name)
SELECT
'CPU_Usage_Percent',
cpu_percent,
'%',
@@SERVERNAME,
DB_NAME()
FROM sys.dm_os_ring_buffers
WHERE ring_buffer_type = 'RING_BUFFER_SCHEDULER_MONITOR';
-- Memory Usage
INSERT INTO performance_metrics (metric_name, metric_value, metric_unit, server_name, database_name)
SELECT
'Memory_Usage_MB',
(physical_memory_in_use_kb / 1024.0),
'MB',
@@SERVERNAME,
DB_NAME()
FROM sys.dm_os_process_memory;
-- Active Connections
INSERT INTO performance_metrics (metric_name, metric_value, metric_unit, server_name, database_name)
SELECT
'Active_Connections',
COUNT(*),
'Count',
@@SERVERNAME,
DB_NAME()
FROM sys.dm_exec_sessions
WHERE is_user_process = 1;
END;
Monitoring Tools: - SQL Server: Extended Events, DMVs, SQL Server Profiler - PostgreSQL: pg_stat_statements, pg_stat_monitor - MySQL: Performance Schema, Information Schema - Cloud: Azure Monitor, AWS CloudWatch, Google Cloud Monitoring
62. How do you troubleshoot database connection issues?
Answer: Database connection troubleshooting follows a systematic approach:
Troubleshooting Steps: 1. Network Connectivity: Check network connectivity and firewall rules 2. Authentication: Verify credentials and permissions 3. Connection Pool: Check connection pool configuration and limits 4. Database Status: Verify database is online and accessible 5. Resource Constraints: Check for resource exhaustion
Implementation Example:
import psycopg2
import logging
from contextlib import contextmanager
from typing import Optional
class DatabaseConnectionTroubleshooter:
def __init__(self, connection_string: str):
self.connection_string = connection_string
self.logger = logging.getLogger(__name__)
def test_connection(self) -> dict:
"""Test database connection and return diagnostic information"""
diagnostics = {
'connection_successful': False,
'error_message': None,
'connection_time_ms': None,
'server_version': None,
'database_name': None
}
try:
import time
start_time = time.time()
with psycopg2.connect(self.connection_string) as conn:
diagnostics['connection_successful'] = True
diagnostics['connection_time_ms'] = (time.time() - start_time) * 1000
diagnostics['server_version'] = conn.server_version
diagnostics['database_name'] = conn.info.dbname
# Test basic query execution
with conn.cursor() as cursor:
cursor.execute("SELECT 1")
cursor.fetchone()
except Exception as e:
diagnostics['error_message'] = str(e)
self.logger.error(f"Connection failed: {e}")
return diagnostics
def check_connection_pool(self) -> dict:
"""Check connection pool status"""
pool_status = {
'active_connections': 0,
'idle_connections': 0,
'total_connections': 0
}
try:
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
# PostgreSQL specific query to check connection status
cursor.execute("""
SELECT
state,
COUNT(*) as connection_count
FROM pg_stat_activity
WHERE datname = current_database()
GROUP BY state
""")
for state, count in cursor.fetchall():
if state == 'active':
pool_status['active_connections'] = count
elif state == 'idle':
pool_status['idle_connections'] = count
pool_status['total_connections'] += count
except Exception as e:
self.logger.error(f"Failed to check connection pool: {e}")
return pool_status
Usage example
troubleshooter = DatabaseConnectionTroubleshooter( "postgresql://user:password@localhost:5432/dbname" )
Test connection
result = troubleshooter.test_connection() print(f"Connection successful: {result['connection_successful']}") print(f"Connection time: {result['connection_time_ms']}ms")
Check pool status
pool_status = troubleshooter.check_connection_pool() print(f"Active connections: {pool_status['active_connections']}")
63. How do you handle database blocking and deadlocks?
Answer: Blocking and deadlock management involves prevention, detection, and resolution strategies.
Prevention Strategies: - Consistent transaction ordering - Proper indexing - Short transaction duration - Appropriate isolation levels
Detection and Resolution:
-- Create a blocking monitoring system
CREATE TABLE blocking_events (
id BIGINT IDENTITY(1,1) PRIMARY KEY,
blocked_session_id INT,
blocking_session_id INT,
blocked_query NVARCHAR(MAX),
blocking_query NVARCHAR(MAX),
wait_time_ms INT,
wait_type NVARCHAR(60),
resource_description NVARCHAR(256),
event_time DATETIME2 DEFAULT GETDATE()
);
-- Stored procedure to detect and log blocking
CREATE PROCEDURE MonitorBlocking
AS
BEGIN
INSERT INTO blocking_events (
blocked_session_id,
blocking_session_id,
blocked_query,
blocking_query,
wait_time_ms,
wait_type,
resource_description
)
SELECT
w.session_id as blocked_session_id,
w.blocking_session_id,
blocked.text as blocked_query,
blocking.text as blocking_query,
w.wait_time,
w.wait_type,
w.resource_description
FROM sys.dm_os_waiting_tasks w
INNER JOIN sys.dm_exec_sessions s ON w.session_id = s.session_id
CROSS APPLY sys.dm_exec_sql_text(s.most_recent_sql_handle) blocked
CROSS APPLY sys.dm_exec_sql_text(
(SELECT most_recent_sql_handle FROM sys.dm_exec_sessions WHERE session_id = w.blocking_session_id)
) blocking
WHERE w.blocking_session_id > 0
AND s.is_user_process = 1;
END;
-- Deadlock detection and resolution
CREATE PROCEDURE HandleDeadlocks
AS
BEGIN
DECLARE @deadlock_victim INT;
-- Get deadlock victim session
SELECT TOP 1 @deadlock_victim = session_id
FROM sys.dm_exec_requests
WHERE status = 'suspended'
AND wait_type = 'LCK_M_X'
ORDER BY wait_time DESC;
IF @deadlock_victim IS NOT NULL
BEGIN
-- Kill the deadlock victim
DECLARE @kill_command NVARCHAR(100) = 'KILL ' + CAST(@deadlock_victim AS NVARCHAR(10));
EXEC sp_executesql @kill_command;
-- Log the deadlock resolution
INSERT INTO blocking_events (
blocked_session_id,
blocking_session_id,
blocked_query,
wait_type,
resource_description
)
VALUES (
@deadlock_victim,
NULL,
'DEADLOCK VICTIM - SESSION KILLED',
'DEADLOCK',
'Automatic deadlock resolution'
);
END
END;
Application-Level Handling:
import psycopg2
from psycopg2.extras import RealDictCursor
import time
import logging
from typing import Optional, Any
class DatabaseTransactionManager:
def __init__(self, connection_string: str):
self.connection_string = connection_string
self.logger = logging.getLogger(__name__)
self.max_retries = 3
self.retry_delay = 1 # seconds
@contextmanager
def transaction_with_retry(self, isolation_level: str = "READ_COMMITTED"):
"""Execute transaction with automatic retry on deadlock"""
retries = 0
while retries < self.max_retries:
try:
with psycopg2.connect(self.connection_string) as conn:
conn.set_isolation_level(getattr(psycopg2.extensions, f"ISOLATION_LEVEL_{isolation_level}"))
with conn.cursor(cursor_factory=RealDictCursor) as cursor:
yield cursor
conn.commit()
return
except psycopg2.extensions.TransactionRollbackError as e:
if "deadlock" in str(e).lower():
retries += 1
self.logger.warning(f"Deadlock detected, retrying ({retries}/{self.max_retries})")
time.sleep(self.retry_delay * retries)
else:
raise
except Exception as e:
self.logger.error(f"Transaction failed: {e}")
raise
raise Exception(f"Transaction failed after {self.max_retries} retries due to deadlocks")
Usage example
transaction_manager = DatabaseTransactionManager("postgresql://user:password@localhost:5432/dbname")
try: with transaction_manager.transaction_with_retry() as cursor: cursor.execute("UPDATE accounts SET balance = balance - 100 WHERE account_id = %s", (123,)) cursor.execute("UPDATE accounts SET balance = balance + 100 WHERE account_id = %s", (456,)) except Exception as e: print(f"Transaction failed: {e}")
64. How do you implement database health checks?
Answer: Database health checks monitor various aspects of database health and provide early warning of issues.
Health Check Components: - Database connectivity - Query performance - Resource utilization - Data integrity - Backup status
Implementation Example:
import psycopg2
import time
import json
from typing import Dict, List, Any
from dataclasses import dataclass
from enum import Enum
class HealthStatus(Enum):
HEALTHY = "healthy"
WARNING = "warning"
CRITICAL = "critical"
UNKNOWN = "unknown"
@dataclass
class HealthCheckResult:
component: str
status: HealthStatus
message: str
metrics: Dict[str, Any]
timestamp: float
class DatabaseHealthChecker:
def __init__(self, connection_string: str):
self.connection_string = connection_string
self.health_checks = [
self.check_connectivity,
self.check_query_performance,
self.check_disk_space,
self.check_connection_pool,
self.check_replication_lag,
self.check_backup_status
]
def check_connectivity(self) -> HealthCheckResult:
"""Check database connectivity"""
start_time = time.time()
try:
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
cursor.execute("SELECT 1")
cursor.fetchone()
response_time = (time.time() - start_time) * 1000
if response_time < 100:
status = HealthStatus.HEALTHY
elif response_time < 500:
status = HealthStatus.WARNING
else:
status = HealthStatus.CRITICAL
return HealthCheckResult(
component="connectivity",
status=status,
message=f"Database connectivity check completed in {response_time:.2f}ms",
metrics={"response_time_ms": response_time},
timestamp=time.time()
)
except Exception as e:
return HealthCheckResult(
component="connectivity",
status=HealthStatus.CRITICAL,
message=f"Database connectivity failed: {str(e)}",
metrics={"error": str(e)},
timestamp=time.time()
)
def check_query_performance(self) -> HealthCheckResult:
"""Check query performance using sample queries"""
try:
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
# Test simple query performance
start_time = time.time()
cursor.execute("SELECT COUNT(*) FROM information_schema.tables")
cursor.fetchone()
simple_query_time = (time.time() - start_time) * 1000
# Test complex query performance
start_time = time.time()
cursor.execute("""
SELECT
schemaname,
tablename,
attname,
n_distinct,
correlation
FROM pg_stats
LIMIT 100
""")
cursor.fetchall()
complex_query_time = (time.time() - start_time) * 1000
# Determine status based on performance
if simple_query_time < 50 and complex_query_time < 200:
status = HealthStatus.HEALTHY
elif simple_query_time < 100 and complex_query_time < 500:
status = HealthStatus.WARNING
else:
status = HealthStatus.CRITICAL
return HealthCheckResult(
component="query_performance",
status=status,
message="Query performance check completed",
metrics={
"simple_query_time_ms": simple_query_time,
"complex_query_time_ms": complex_query_time
},
timestamp=time.time()
)
except Exception as e:
return HealthCheckResult(
component="query_performance",
status=HealthStatus.CRITICAL,
message=f"Query performance check failed: {str(e)}",
metrics={"error": str(e)},
timestamp=time.time()
)
def check_disk_space(self) -> HealthCheckResult:
"""Check available disk space"""
try:
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
cursor.execute("""
SELECT
pg_database_size(current_database()) as db_size,
pg_size_pretty(pg_database_size(current_database())) as db_size_pretty
""")
result = cursor.fetchone()
db_size_bytes = result[0]
db_size_pretty = result[1]
# This is a simplified check - in production you'd check actual disk space
# For demonstration, we'll use database size as a proxy
if db_size_bytes < 1024 * 1024 * 1024: # 1GB
status = HealthStatus.HEALTHY
elif db_size_bytes < 5 * 1024 * 1024 * 1024: # 5GB
status = HealthStatus.WARNING
else:
status = HealthStatus.CRITICAL
return HealthCheckResult(
component="disk_space",
status=status,
message=f"Database size: {db_size_pretty}",
metrics={"database_size_bytes": db_size_bytes},
timestamp=time.time()
)
except Exception as e:
return HealthCheckResult(
component="disk_space",
status=HealthStatus.CRITICAL,
message=f"Disk space check failed: {str(e)}",
metrics={"error": str(e)},
timestamp=time.time()
)
def run_all_checks(self) -> Dict[str, HealthCheckResult]:
"""Run all health checks"""
results = {}
for check in self.health_checks:
try:
result = check()
results[result.component] = result
except Exception as e:
results[check.__name__] = HealthCheckResult(
component=check.__name__,
status=HealthStatus.UNKNOWN,
message=f"Health check failed: {str(e)}",
metrics={"error": str(e)},
timestamp=time.time()
)
return results
def get_overall_status(self, results: Dict[str, HealthCheckResult]) -> HealthStatus:
"""Determine overall health status"""
if any(result.status == HealthStatus.CRITICAL for result in results.values()):
return HealthStatus.CRITICAL
elif any(result.status == HealthStatus.WARNING for result in results.values()):
return HealthStatus.WARNING
elif all(result.status == HealthStatus.HEALTHY for result in results.values()):
return HealthStatus.HEALTHY
else:
return HealthStatus.UNKNOWN
Usage example
health_checker = DatabaseHealthChecker("postgresql://user:password@localhost:5432/dbname") results = health_checker.run_all_checks() overall_status = health_checker.get_overall_status(results)
print(f"Overall database health: {overall_status.value}") for component, result in results.items(): print(f"{component}: {result.status.value} - {result.message}")
Database Operations & Maintenance
65. How do you handle database corruption scenarios?
Answer: Database corruption requires immediate detection, isolation, and recovery procedures.
Key Strategies: - Prevention: Regular backups, checksums, write-ahead logging - Detection: Automated monitoring, integrity checks, error logging - Recovery: Point-in-time recovery, transaction log replay, data reconstruction
Coding Example - Corruption Detection Script:
import psycopg2
import logging
from datetime import datetime
class DatabaseCorruptionHandler:
def __init__(self, connection_string):
self.connection_string = connection_string
self.logger = logging.getLogger(__name__)
def check_database_integrity(self):
"""Perform comprehensive database integrity checks"""
try:
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
# Check for corrupted pages
cursor.execute("""
SELECT schemaname, tablename, attname,
pg_stat_get_tuples_returned(c.oid) as rows_returned
FROM pg_class c
JOIN pg_attribute a ON a.attrelid = c.oid
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relkind = 'r'
AND pg_stat_get_tuples_returned(c.oid) < 0
""")
corrupted_tables = cursor.fetchall()
if corrupted_tables:
self.logger.error(f"Corruption detected in tables: {corrupted_tables}")
return False
# Check for orphaned files
cursor.execute("""
SELECT relname, relfilenode
FROM pg_class c
LEFT JOIN pg_stat_file('base/' || c.relfilenode) f ON true
WHERE c.relkind = 'r' AND f.size IS NULL
""")
orphaned_files = cursor.fetchall()
if orphaned_files:
self.logger.error(f"Orphaned files detected: {orphaned_files}")
return False
return True
except Exception as e:
self.logger.error(f"Integrity check failed: {str(e)}")
return False
def recover_corrupted_table(self, table_name, schema_name='public'):
"""Attempt to recover a corrupted table"""
try:
with psycopg2.connect(self.connection_string) as conn:
conn.autocommit = False
# Create backup of corrupted table
backup_table = f"{table_name}_corrupted_backup_{datetime.now().strftime('%Y%m%d_%H%M%S')}"
with conn.cursor() as cursor:
# Attempt to create backup
cursor.execute(f"""
CREATE TABLE {backup_table} AS
SELECT * FROM {schema_name}.{table_name}
WHERE 1=0
""")
# Try to copy data in chunks
cursor.execute(f"SELECT COUNT(*) FROM {schema_name}.{table_name}")
total_rows = cursor.fetchone()[0]
chunk_size = 1000
for offset in range(0, total_rows, chunk_size):
try:
cursor.execute(f"""
INSERT INTO {backup_table}
SELECT * FROM {schema_name}.{table_name}
LIMIT {chunk_size} OFFSET {offset}
""")
except Exception as chunk_error:
self.logger.warning(f"Failed to copy chunk {offset}: {chunk_error}")
continue
conn.commit()
self.logger.info(f"Recovery completed for table {table_name}")
except Exception as e:
self.logger.error(f"Recovery failed: {str(e)}")
conn.rollback()
raise
66. How do you monitor database space usage?
Answer: Implement comprehensive monitoring for storage, indexes, and growth trends.
Coding Example - Space Monitoring System:
import psycopg2
import pandas as pd
from datetime import datetime, timedelta
import smtplib
from email.mime.text import MIMEText
class DatabaseSpaceMonitor:
def __init__(self, connection_string, alert_threshold=80):
self.connection_string = connection_string
self.alert_threshold = alert_threshold
def get_database_size_info(self):
"""Get comprehensive database size information"""
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
# Database sizes
cursor.execute("""
SELECT
d.datname as database_name,
pg_size_pretty(pg_database_size(d.datname)) as size_pretty,
pg_database_size(d.datname) as size_bytes,
pg_database_size(d.datname) / 1024 / 1024 / 1024.0 as size_gb
FROM pg_database d
WHERE d.datname NOT IN ('template0', 'template1')
ORDER BY pg_database_size(d.datname) DESC
""")
db_sizes = cursor.fetchall()
# Table sizes
cursor.execute("""
SELECT
schemaname,
tablename,
pg_size_pretty(pg_total_relation_size(schemaname||'.'||tablename)) as size_pretty,
pg_total_relation_size(schemaname||'.'||tablename) as size_bytes,
pg_total_relation_size(schemaname||'.'||tablename) / 1024 / 1024.0 as size_mb
FROM pg_tables
WHERE schemaname NOT IN ('information_schema', 'pg_catalog')
ORDER BY pg_total_relation_size(schemaname||'.'||tablename) DESC
LIMIT 20
""")
table_sizes = cursor.fetchall()
# Index sizes
cursor.execute("""
SELECT
schemaname,
tablename,
indexname,
pg_size_pretty(pg_relation_size(schemaname||'.'||indexname)) as size_pretty,
pg_relation_size(schemaname||'.'||indexname) as size_bytes
FROM pg_indexes
WHERE schemaname NOT IN ('information_schema', 'pg_catalog')
ORDER BY pg_relation_size(schemaname||'.'||indexname) DESC
LIMIT 20
""")
index_sizes = cursor.fetchall()
return {
'database_sizes': db_sizes,
'table_sizes': table_sizes,
'index_sizes': index_sizes
}
def check_disk_space(self):
"""Check available disk space"""
import shutil
total, used, free = shutil.disk_usage("/")
usage_percent = (used / total) * 100
return {
'total_gb': total / (1024**3),
'used_gb': used / (1024**3),
'free_gb': free / (1024**3),
'usage_percent': usage_percent
}
def generate_space_report(self):
"""Generate comprehensive space usage report"""
space_info = self.get_database_size_info()
disk_info = self.check_disk_space()
report = f"""
Database Space Usage Report - {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}
DISK USAGE:
Total: {disk_info['total_gb']:.2f} GB
Used: {disk_info['used_gb']:.2f} GB ({disk_info['usage_percent']:.1f}%)
Free: {disk_info['free_gb']:.2f} GB
TOP DATABASES BY SIZE:
"""
for db in space_info['database_sizes'][:10]:
report += f"{db[0]}: {db[1]} ({db[3]:.2f} GB)\n"
report += "\nTOP TABLES BY SIZE:\n"
for table in space_info['table_sizes'][:10]:
report += f"{table[0]}.{table[1]}: {table[2]} ({table[4]:.2f} MB)\n"
return report
def send_alert_if_needed(self, report):
"""Send alert if space usage exceeds threshold"""
disk_info = self.check_disk_space()
if disk_info['usage_percent'] > self.alert_threshold:
# Send email alert
msg = MIMEText(f"Database space usage alert!\n\n{report}")
msg['Subject'] = f"Database Space Alert - {disk_info['usage_percent']:.1f}% used"
msg['From'] = "db-monitor@company.com"
msg['To'] = "dba-team@company.com"
# Send email (configure SMTP settings)
# smtp_server.send_message(msg)
return True
return False
67. How do you troubleshoot database timeout issues?
Answer: Systematic approach to identify and resolve timeout root causes.
Coding Example - Timeout Diagnostics:
import psycopg2
import time
import threading
from collections import defaultdict
class DatabaseTimeoutDiagnostics:
def __init__(self, connection_string):
self.connection_string = connection_string
self.timeout_threshold = 30 # seconds
def monitor_active_queries(self, duration=60):
"""Monitor active queries for timeout patterns"""
start_time = time.time()
query_stats = defaultdict(list)
while time.time() - start_time < duration:
try:
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
cursor.execute("""
SELECT
pid,
usename,
application_name,
client_addr,
state,
query_start,
state_change,
query,
wait_event_type,
wait_event
FROM pg_stat_activity
WHERE state = 'active'
AND query NOT LIKE '%pg_stat_activity%'
""")
active_queries = cursor.fetchall()
for query in active_queries:
query_duration = time.time() - query[5].timestamp()
query_stats[query[0]].append({
'duration': query_duration,
'query': query[7][:100], # Truncate long queries
'user': query[1],
'application': query[2],
'wait_event': query[8]
})
if query_duration > self.timeout_threshold:
print(f"LONG RUNNING QUERY DETECTED: PID {query[0]}, Duration: {query_duration:.2f}s")
time.sleep(5) # Check every 5 seconds
except Exception as e:
print(f"Monitoring error: {e}")
time.sleep(5)
return query_stats
def analyze_query_performance(self, query_text):
"""Analyze specific query performance"""
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
# Enable query timing
cursor.execute("SET log_statement = 'all'")
cursor.execute("SET log_min_duration_statement = 1000") # Log queries > 1s
# Get execution plan
cursor.execute(f"EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) {query_text}")
plan = cursor.fetchone()[0]
# Get table statistics
cursor.execute("""
SELECT schemaname, tablename, attname, n_distinct, correlation
FROM pg_stats
WHERE schemaname NOT IN ('information_schema', 'pg_catalog')
ORDER BY schemaname, tablename, attname
""")
stats = cursor.fetchall()
return {
'execution_plan': plan,
'table_statistics': stats
}
def identify_bottlenecks(self):
"""Identify common timeout bottlenecks"""
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
# Check for locks
cursor.execute("""
SELECT
l.pid,
l.mode,
l.granted,
a.usename,
a.query,
a.query_start
FROM pg_locks l
JOIN pg_stat_activity a ON l.pid = a.pid
WHERE NOT l.granted
""")
locks = cursor.fetchall()
# Check for long-running transactions
cursor.execute("""
SELECT
pid,
usename,
application_name,
xact_start,
query_start,
state,
query
FROM pg_stat_activity
WHERE xact_start < NOW() - INTERVAL '5 minutes'
AND state IN ('active', 'idle in transaction')
""")
long_transactions = cursor.fetchall()
# Check connection count
cursor.execute("""
SELECT
usename,
application_name,
COUNT(*) as connection_count
FROM pg_stat_activity
GROUP BY usename, application_name
ORDER BY connection_count DESC
""")
connections = cursor.fetchall()
return {
'blocked_queries': locks,
'long_transactions': long_transactions,
'connection_distribution': connections
}
def optimize_timeout_settings(self):
"""Suggest timeout optimization settings"""
recommendations = {
'statement_timeout': '30s', # Prevent runaway queries
'lock_timeout': '10s', # Prevent indefinite lock waits
'idle_in_transaction_session_timeout': '5min', # Clean up idle transactions
'max_connections': '100', # Limit concurrent connections
'shared_buffers': '25% of RAM', # Optimize memory usage
'effective_cache_size': '75% of RAM' # Help query planner
}
return recommendations
68. How do you implement database performance baselines?
Answer: Establish comprehensive performance metrics and monitoring systems.
Coding Example - Performance Baseline System:
import psycopg2
import json
import time
from datetime import datetime, timedelta
import statistics
class DatabasePerformanceBaseline:
def __init__(self, connection_string):
self.connection_string = connection_string
self.baseline_data = {}
def collect_performance_metrics(self):
"""Collect comprehensive performance metrics"""
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
# Query execution statistics
cursor.execute("""
SELECT
query,
calls,
total_time,
mean_time,
stddev_time,
min_time,
max_time
FROM pg_stat_statements
ORDER BY total_time DESC
LIMIT 50
""")
query_stats = cursor.fetchall()
# Table access patterns
cursor.execute("""
SELECT
schemaname,
tablename,
seq_scan,
seq_tup_read,
idx_scan,
idx_tup_fetch,
n_tup_ins,
n_tup_upd,
n_tup_del,
n_live_tup,
n_dead_tup
FROM pg_stat_user_tables
ORDER BY n_live_tup DESC
""")
table_stats = cursor.fetchall()
# Index usage statistics
cursor.execute("""
SELECT
schemaname,
tablename,
indexname,
idx_scan,
idx_tup_read,
idx_tup_fetch
FROM pg_stat_user_indexes
ORDER BY idx_scan DESC
""")
index_stats = cursor.fetchall()
# Connection and transaction statistics
cursor.execute("""
SELECT
datname,
numbackends,
xact_commit,
xact_rollback,
blks_read,
blks_hit,
tup_returned,
tup_fetched,
tup_inserted,
tup_updated,
tup_deleted
FROM pg_stat_database
WHERE datname = current_database()
""")
db_stats = cursor.fetchone()
return {
'timestamp': datetime.now().isoformat(),
'query_statistics': query_stats,
'table_statistics': table_stats,
'index_statistics': index_stats,
'database_statistics': db_stats
}
def establish_baseline(self, collection_duration=3600, interval=300):
"""Establish performance baseline over specified duration"""
baseline_metrics = []
start_time = time.time()
while time.time() - start_time < collection_duration:
metrics = self.collect_performance_metrics()
baseline_metrics.append(metrics)
time.sleep(interval)
# Calculate baseline statistics
baseline = self.calculate_baseline_statistics(baseline_metrics)
# Save baseline to file
with open('performance_baseline.json', 'w') as f:
json.dump(baseline, f, indent=2)
return baseline
def calculate_baseline_statistics(self, metrics_list):
"""Calculate statistical baseline from collected metrics"""
baseline = {
'created_at': datetime.now().isoformat(),
'sample_count': len(metrics_list),
'query_performance': {},
'table_performance': {},
'system_performance': {}
}
# Analyze query performance
query_times = defaultdict(list)
for metrics in metrics_list:
for query_stat in metrics['query_statistics']:
query_hash = hash(query_stat[0]) # Simple hash of query
query_times[query_hash].append(query_stat[3]) # mean_time
for query_hash, times in query_times.items():
baseline['query_performance'][str(query_hash)] = {
'mean': statistics.mean(times),
'median': statistics.median(times),
'std_dev': statistics.stdev(times) if len(times) > 1 else 0,
'min': min(times),
'max': max(times),
'count': len(times)
}
# Analyze table performance
table_access = defaultdict(lambda: {
'seq_scans': [], 'idx_scans': [], 'live_tuples': []
})
for metrics in metrics_list:
for table_stat in metrics['table_statistics']:
table_key = f"{table_stat[0]}.{table_stat[1]}"
table_access[table_key]['seq_scans'].append(table_stat[2])
table_access[table_key]['idx_scans'].append(table_stat[4])
table_access[table_key]['live_tuples'].append(table_stat[9])
for table_key, stats in table_access.items():
baseline['table_performance'][table_key] = {
'avg_seq_scans': statistics.mean(stats['seq_scans']),
'avg_idx_scans': statistics.mean(stats['idx_scans']),
'avg_live_tuples': statistics.mean(stats['live_tuples'])
}
return baseline
def compare_with_baseline(self, current_metrics):
"""Compare current metrics with established baseline"""
try:
with open('performance_baseline.json', 'r') as f:
baseline = json.load(f)
except FileNotFoundError:
return {"error": "No baseline found"}
deviations = {
'query_deviations': [],
'table_deviations': [],
'overall_health': 'good'
}
# Check query performance deviations
for query_hash, baseline_stats in baseline['query_performance'].items():
if query_hash in current_metrics['query_statistics']:
current_time = current_metrics['query_statistics'][query_hash][3]
baseline_mean = baseline_stats['mean']
if current_time > baseline_mean * 2: # 2x slower
deviations['query_deviations'].append({
'query_hash': query_hash,
'baseline_time': baseline_mean,
'current_time': current_time,
'deviation_percent': ((current_time - baseline_mean) / baseline_mean) * 100
})
# Check table performance deviations
for table_key, baseline_stats in baseline['table_performance'].items():
if table_key in current_metrics['table_statistics']:
current_seq_scans = current_metrics['table_statistics'][table_key][2]
baseline_seq_scans = baseline_stats['avg_seq_scans']
if current_seq_scans > baseline_seq_scans * 1.5: # 50% more seq scans
deviations['table_deviations'].append({
'table': table_key,
'baseline_seq_scans': baseline_seq_scans,
'current_seq_scans': current_seq_scans,
'deviation_percent': ((current_seq_scans - baseline_seq_scans) / baseline_seq_scans) * 100
})
if deviations['query_deviations'] or deviations['table_deviations']:
deviations['overall_health'] = 'degraded'
return deviations
69. How do you handle database upgrade scenarios?
Answer: Systematic approach to planning, testing, and executing database upgrades.
Coding Example - Database Upgrade Framework:
import psycopg2
import subprocess
import yaml
import logging
from datetime import datetime
import hashlib
class DatabaseUpgradeManager:
def __init__(self, connection_string, upgrade_config_path):
self.connection_string = connection_string
self.config = self.load_upgrade_config(upgrade_config_path)
self.logger = logging.getLogger(__name__)
def load_upgrade_config(self, config_path):
"""Load upgrade configuration from YAML file"""
with open(config_path, 'r') as f:
return yaml.safe_load(f)
def get_current_version(self):
"""Get current database version"""
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
cursor.execute("""
CREATE TABLE IF NOT EXISTS db_version (
version VARCHAR(50) PRIMARY KEY,
applied_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
checksum VARCHAR(64),
description TEXT
)
""")
cursor.execute("SELECT version FROM db_version ORDER BY applied_at DESC LIMIT 1")
result = cursor.fetchone()
return result[0] if result else '0.0.0'
def calculate_checksum(self, sql_file_path):
"""Calculate checksum of SQL file"""
with open(sql_file_path, 'rb') as f:
return hashlib.sha256(f.read()).hexdigest()
def backup_database(self, backup_path):
"""Create database backup before upgrade"""
timestamp = datetime.now().strftime('%Y%m%d_%H%M%S')
backup_file = f"{backup_path}/backup_{timestamp}.sql"
try:
# Extract connection details
conn_parts = self.connection_string.split()
db_name = None
for part in conn_parts:
if part.startswith('dbname='):
db_name = part.split('=')[1]
break
if not db_name:
raise ValueError("Database name not found in connection string")
# Create backup using pg_dump
cmd = [
'pg_dump',
'-h', 'localhost', # Configure as needed
'-U', 'postgres', # Configure as needed
'-d', db_name,
'-f', backup_file,
'--verbose'
]
result = subprocess.run(cmd, capture_output=True, text=True)
if result.returncode == 0:
self.logger.info(f"Backup created successfully: {backup_file}")
return backup_file
else:
raise Exception(f"Backup failed: {result.stderr}")
except Exception as e:
self.logger.error(f"Backup failed: {str(e)}")
raise
def apply_upgrade_script(self, script_path, version, description):
"""Apply individual upgrade script"""
with psycopg2.connect(self.connection_string) as conn:
conn.autocommit = False
try:
with conn.cursor() as cursor:
# Read and execute SQL script
with open(script_path, 'r') as f:
sql_content = f.read()
# Split into individual statements
statements = [stmt.strip() for stmt in sql_content.split(';') if stmt.strip()]
for statement in statements:
if statement:
cursor.execute(statement)
# Record upgrade in version table
checksum = self.calculate_checksum(script_path)
cursor.execute("""
INSERT INTO db_version (version, checksum, description)
VALUES (%s, %s, %s)
""", (version, checksum, description))
conn.commit()
self.logger.info(f"Successfully applied upgrade {version}")
except Exception as e:
conn.rollback()
self.logger.error(f"Upgrade {version} failed: {str(e)}")
raise
def validate_upgrade(self, version):
"""Validate upgrade was applied correctly"""
with psycopg2.connect(self.connection_string) as conn:
with conn.cursor() as cursor:
# Check version was recorded
cursor.execute("SELECT version FROM db_version WHERE version = %s", (version,))
if not cursor.fetchone():
return False, f"Version {version} not found in db_version table"
# Run validation queries from config
if 'validations' in self.config and version in self.config['validations']:
for validation in self.config['validations'][version]:
cursor.execute(validation['query'])
result = cursor.fetchone()[0]
if not validation['expected_result'](result):
return False, f"Validation failed: {validation['description']}"
return True, "Validation passed"
def rollback_upgrade(self, version):
"""Rollback specific upgrade version"""
rollback_script = self.config['rollbacks'].get(version)
if not rollback_script:
raise ValueError(f"No rollback script found for version {version}")
with psycopg2.connect(self.connection_string) as conn:
conn.autocommit = False
try:
with conn.cursor() as cursor:
# Apply rollback script
with open(rollback_script, 'r') as f:
sql_content = f.read()
statements = [stmt.strip() for stmt in sql_content.split(';') if stmt.strip()]
for statement in statements:
if statement:
cursor.execute(statement)
# Remove version record
cursor.execute("DELETE FROM db_version WHERE version = %s", (version,))
conn.commit()
self.logger.info(f"Successfully rolled back version {version}")
except Exception as e:
conn.rollback()
self.logger.error(f"Rollback failed: {str(e)}")
raise
def execute_upgrade_plan(self, target_version):
"""Execute complete upgrade plan"""
current_version = self.get_current_version()
# Create backup
backup_file = self.backup_database(self.config['backup_path'])
try:
# Get upgrade path
upgrade_path = self.calculate_upgrade_path(current_version, target_version)
for upgrade_step in upgrade_path:
self.logger.info(f"Applying upgrade: {upgrade_step['version']}")
# Apply upgrade
self.apply_upgrade_script(
upgrade_step['script_path'],
upgrade_step['version'],
upgrade_step['description']
)
# Validate upgrade
success, message = self.validate_upgrade(upgrade_step['version'])
if not success:
raise Exception(f"Validation failed: {message}")
self.logger.info(f"Upgrade {upgrade_step['version']} completed successfully")
self.logger.info(f"Database upgraded from {current_version} to {target_version}")
except Exception as e:
self.logger.error(f"Upgrade failed: {str(e)}")
self.logger.info("Consider restoring from backup if needed")
raise
def calculate_upgrade_path(self, current_version, target_version):
"""Calculate upgrade path between versions"""
# This would implement version comparison logic
# For simplicity, assuming sequential upgrades
upgrades = []
for version_info in self.config['upgrades']:
if self.version_greater_than(version_info['version'], current_version) and \
not self.version_greater_than(version_info['version'], target_version):
upgrades.append(version_info)
return sorted(upgrades, key=lambda x: x['version'])
def version_greater_than(self, version1, version2):
"""Compare version strings"""
v1_parts = [int(x) for x in version1.split('.')]
v2_parts = [int(x) for x in version2.split('.')]
for i in range(max(len(v1_parts), len(v2_parts))):
v1_part = v1_parts[i] if i < len(v1_parts) else 0
v2_part = v2_parts[i] if i < len(v2_parts) else 0
if v1_part > v2_part:
return True
elif v1_part < v2_part:
return False
return False
70. How do you implement database disaster recovery testing?
Answer: Comprehensive testing framework for disaster recovery procedures.
Coding Example - Disaster Recovery Testing Framework:
import psycopg2
import subprocess
import time
import json
from datetime import datetime, timedelta
import threading
class DisasterRecoveryTester:
def __init__(self, primary_connection, backup_connection, test_config):
self.primary_connection = primary_connection
self.backup_connection = backup_connection
self.test_config = test_config
self.test_results = []
def test_backup_integrity(self):
"""Test backup file integrity and restorability"""
test_result = {
'test_name': 'backup_integrity',
'timestamp': datetime.now().isoformat(),
'status': 'unknown',
'details': {}
}
try:
# Create test backup
backup_file = self.create_test_backup()
# Verify backup file exists and has content
import os
if not os.path.exists(backup_file):
raise Exception("Backup file not created")
file_size = os.path.getsize(backup_file)
if file_size == 0:
raise Exception("Backup file is empty")
test_result['details']['backup_size_mb'] = file_size / (1024 * 1024)
# Test backup restore to temporary database
restore_success = self.test_backup_restore(backup_file)
test_result['details']['restore_success'] = restore_success
if restore_success:
test_result['status'] = 'passed'
else:
test_result['status'] = 'failed'
except Exception as e:
test_result['status'] = 'failed'
test_result['details']['error'] = str(e)
self.test_results.append(test_result)
return test_result
def test_failover_scenario(self):
"""Test database failover to backup system"""
test_result = {
'test_name': 'failover_scenario',
'timestamp': datetime.now().isoformat(),
'status': 'unknown',
'details': {}
}
try:
# Simulate primary database failure
self.simulate_primary_failure()
# Measure failover time
start_time = time.time()
# Perform failover
failover_success = self.perform_failover()
failover_time = time.time() - start_time
test_result['details']['failover_time_seconds'] = failover_time
if failover_success:
# Test application connectivity to backup
connectivity_success = self.test_backup_connectivity()
test_result['details']['connectivity_success'] = connectivity_success
# Test data integrity after failover
data_integrity = self.test_data_integrity()
test_result['details']['data_integrity'] = data_integrity
if connectivity_success and data_integrity:
test_result['status'] = 'passed'
else:
test_result['status'] = 'failed'
else:
test_result['status'] = 'failed'
# Restore primary database
self.restore_primary_database()
except Exception as e:
test_result['status'] = 'failed'
test_result['details']['error'] = str(e)
self.test_results.append(test_result)
return test_result
def test_point_in_time_recovery(self):
"""Test point-in-time recovery capabilities"""
test_result = {
'test_name': 'point_in_time_recovery',
'timestamp': datetime.now().isoformat(),
'status': 'unknown',
'details': {}
}
try:
# Create test data
test_data_id = self.create_test_data()
# Record recovery point
recovery_point = datetime.now()
test_result['details']['recovery_point'] = recovery_point.isoformat()
# Simulate data corruption
self.simulate_data_corruption(test_data_id)
# Perform point-in-time recovery
recovery_success = self.perform_point_in_time_recovery(recovery_point)
test_result['details']['recovery_success'] = recovery_success
if recovery_success:
# Verify data was restored correctly
data_restored = self.verify_test_data_restored(test_data_id)
test_result['details']['data_restored'] = data_restored
if data_restored:
test_result['status'] = 'passed'
else:
test_result['status'] = 'failed'
else:
test_result['status'] = 'failed'
# Clean up test data
self.cleanup_test_data(test_data_id)
except Exception as e:
test_result['status'] = 'failed'
test_result['details']['error'] = str(e)
self.test_results.append(test_result)
return test_result
def test_replication_lag(self):
"""Test replication lag and consistency"""
test_result = {
'test_name': 'replication_lag',
'timestamp': datetime.now().isoformat(),
'status': 'unknown',
'details': {}
}
try:
# Measure initial replication lag
initial_lag = self.measure_replication_lag()
test_result['details']['initial_lag_seconds'] = initial_lag
# Create test transaction
test_data_id = self.create_test_data()
# Measure replication lag after transaction
post_tx_lag = self.measure_replication_lag()
test_result['details']['post_transaction_lag_seconds'] = post_tx_lag
# Verify data consistency
consistency_check = self.verify_replication_consistency(test_data_id)
test_result['details']['consistency_check'] = consistency_check
# Check if lag is within acceptable limits
max_acceptable_lag = self.test_config.get('max_replication_lag_seconds', 30)
if post_tx_lag <= max_acceptable_lag and consistency_check:
test_result['status'] = 'passed'
else:
test_result['status'] = 'failed'
# Clean up
self.cleanup_test_data(test_data_id)
except Exception as e:
test_result['status'] = 'failed'
test_result['details']['error'] = str(e)
self.test_results.append(test_result)
return test_result
def test_automated_recovery_procedures(self):
"""Test automated recovery procedures"""
test_result = {
'test_name': 'automated_recovery',
'timestamp': datetime.now().isoformat(),
'status': 'unknown',
'details': {}
}
try:
# Test automated backup verification
backup_verification = self.test_automated_backup_verification()
test_result['details']['backup_verification'] = backup_verification
# Test automated failover detection
failover_detection = self.test_failover_detection()
test_result['details']['failover_detection'] = failover_detection
# Test automated recovery scripts
recovery_scripts = self.test_recovery_scripts()
test_result['details']['recovery_scripts'] = recovery_scripts
# Test monitoring and alerting
monitoring_alerts = self.test_monitoring_alerts()
test_result['details']['monitoring_alerts'] = monitoring_alerts
if all([backup_verification, failover_detection, recovery_scripts, monitoring_alerts]):
test_result['status'] = 'passed'
else:
test_result['status'] = 'failed'
except Exception as e:
test_result['status'] = 'failed'
test_result['details']['error'] = str(e)
self.test_results.append(test_result)
return test_result
def run_comprehensive_dr_test(self):
"""Run comprehensive disaster recovery test suite"""
test_suite = [
self.test_backup_integrity,
self.test_failover_scenario,
self.test_point_in_time_recovery,
self.test_replication_lag,
self.test_automated_recovery_procedures
]
print("Starting Disaster Recovery Test Suite...")
for test_func in test_suite:
print(f"Running {test_func.__name__}...")
result = test_func()
if result['status'] == 'passed':
print(f"✓ {test_func.__name__} PASSED")
else:
print(f"✗ {test_func.__name__} FAILED")
print(f" Error: {result['details'].get('error', 'Unknown error')}")
# Generate test report
self.generate_test_report()
return self.test_results
def generate_test_report(self):
"""Generate comprehensive test report"""
report = {
'test_suite': 'Disaster Recovery',
'execution_time': datetime.now().isoformat(),
'summary': {
'total_tests': len(self.test_results),
'passed': len([r for r in self.test_results if r['status'] == 'passed']),
'failed': len([r for r in self.test_results if r['status'] == 'failed'])
},
'test_results': self.test_results,
'recommendations': self.generate_recommendations()
}
# Save report
with open('dr_test_report.json', 'w') as f:
json.dump(report, f, indent=2)
# Print summary
print(f"\nTest Summary:")
print(f"Total Tests: {report['summary']['total_tests']}")
print(f"Passed: {report['summary']['passed']}")
print(f"Failed: {report['summary']['failed']}")
return report
def generate_recommendations(self):
"""Generate recommendations based on test results"""
recommendations = []
failed_tests = [r for r in self.test_results if r['status'] == 'failed']
for test in failed_tests:
if test['test_name'] == 'backup_integrity':
recommendations.append("Review backup procedures and verify backup file integrity")
elif test['test_name'] == 'failover_scenario':
recommendations.append("Improve failover procedures and reduce failover time")
elif test['test_name'] == 'point_in_time_recovery':
recommendations.append("Enhance point-in-time recovery procedures")
elif test['test_name'] == 'replication_lag':
recommendations.append("Optimize replication configuration to reduce lag")
elif test['test_name'] == 'automated_recovery':
recommendations.append("Improve automated recovery procedures and monitoring")
return recommendations
71. How do you implement database change data capture?
Answer: Change Data Capture (CDC) tracks and captures changes made to database tables, enabling real-time data replication, audit trails, and event-driven architectures.
Implementation Approaches: - Database-level CDC: Using SQL Server CDC, Oracle GoldenGate - Application-level CDC: Using triggers, audit tables, or event sourcing - Log-based CDC: Reading database transaction logs
// Database-level CDC with SQL Server
public class SqlServerCdcService
{
private readonly string _connectionString;
public SqlServerCdcService(string connectionString)
{
_connectionString = connectionString;
}
public async Task EnableCdcForTable(string databaseName, string schemaName, string tableName)
{
using var connection = new SqlConnection(_connectionString);
await connection.OpenAsync();
// Enable CDC for database
var enableDbCdc = $"EXEC sys.sp_cdc_enable_db";
using var cmd1 = new SqlCommand(enableDbCdc, connection);
await cmd1.ExecuteNonQueryAsync();
// Enable CDC for specific table
var enableTableCdc = $"EXEC sys.sp_cdc_enable_table @source_schema = '{schemaName}', @source_name = '{tableName}', @role_name = NULL";
using var cmd2 = new SqlCommand(enableTableCdc, connection);
await cmd2.ExecuteNonQueryAsync();
}
public async Task<List<CdcChange>> GetChanges(string captureInstance, DateTime fromLsn, DateTime toLsn)
{
using var connection = new SqlConnection(_connectionString);
await connection.OpenAsync();
var query = @"
SELECT
__$operation,
__$seqval,
__$update_mask,
*
FROM cdc.fn_cdc_get_all_changes_dbo_" + captureInstance +
$"(@from_lsn, @to_lsn, N'all')";
using var cmd = new SqlCommand(query, connection);
cmd.Parameters.AddWithValue("@from_lsn", fromLsn);
cmd.Parameters.AddWithValue("@to_lsn", toLsn);
var changes = new List<CdcChange>();
using var reader = await cmd.ExecuteReaderAsync();
while (await reader.ReadAsync())
{
changes.Add(new CdcChange
{
Operation = (CdcOperation)reader.GetInt32("__$operation"),
SequenceValue = reader.GetInt64("__$seqval"),
UpdateMask = reader.GetSqlBinary("__$update_mask"),
Data = ReadRowData(reader)
});
}
return changes;
}
}
public enum CdcOperation
{
Delete = 1,
Insert = 2,
Update = 3,
UpdateWithBeforeImage = 4
}
public class CdcChange
{
public CdcOperation Operation { get; set; }
public long SequenceValue { get; set; }
public SqlBinary UpdateMask { get; set; }
public Dictionary<string, object> Data { get; set; }
}
72. How do you handle database integration with message queues?
Answer: Database integration with message queues enables reliable, asynchronous processing, event-driven architectures, and distributed system communication.
// Database-Queue Integration Service
public class DatabaseQueueIntegrationService
{
private readonly IMessageQueue _messageQueue;
private readonly IDbContext _dbContext;
private readonly ILogger<DatabaseQueueIntegrationService> _logger;
public DatabaseQueueIntegrationService(
IMessageQueue messageQueue,
IDbContext dbContext,
ILogger<DatabaseQueueIntegrationService> logger)
{
_messageQueue = messageQueue;
_dbContext = dbContext;
_logger = logger;
}
// Outbox Pattern Implementation
public async Task<Guid> CreateOrderWithOutbox(Order order)
{
using var transaction = await _dbContext.Database.BeginTransactionAsync();
try
{
// Save order to database
_dbContext.Orders.Add(order);
await _dbContext.SaveChangesAsync();
// Create outbox message
var outboxMessage = new OutboxMessage
{
Id = Guid.NewGuid(),
Type = "OrderCreated",
Data = JsonSerializer.Serialize(order),
CreatedAt = DateTime.UtcNow,
ProcessedAt = null
};
_dbContext.OutboxMessages.Add(outboxMessage);
await _dbContext.SaveChangesAsync();
await transaction.CommitAsync();
return order.Id;
}
catch
{
await transaction.RollbackAsync();
throw;
}
}
// Background service to process outbox
public async Task ProcessOutboxMessages()
{
var unprocessedMessages = await _dbContext.OutboxMessages
.Where(m => m.ProcessedAt == null)
.OrderBy(m => m.CreatedAt)
.Take(100)
.ToListAsync();
foreach (var message in unprocessedMessages)
{
try
{
await _messageQueue.PublishAsync(message.Type, message.Data);
message.ProcessedAt = DateTime.UtcNow;
await _dbContext.SaveChangesAsync();
}
catch (Exception ex)
{
_logger.LogError(ex, "Failed to process outbox message {MessageId}", message.Id);
}
}
}
// Event sourcing with queue integration
public async Task AppendEvent(DomainEvent domainEvent)
{
using var transaction = await _dbContext.Database.BeginTransactionAsync();
try
{
// Store event in event store
var eventStore = new EventStore
{
Id = Guid.NewGuid(),
AggregateId = domainEvent.AggregateId,
Version = domainEvent.Version,
EventType = domainEvent.GetType().Name,
EventData = JsonSerializer.Serialize(domainEvent),
Timestamp = DateTime.UtcNow
};
_dbContext.EventStore.Add(eventStore);
await _dbContext.SaveChangesAsync();
// Publish to queue for event handlers
await _messageQueue.PublishAsync("DomainEvent", JsonSerializer.Serialize(domainEvent));
await transaction.CommitAsync();
}
catch
{
await transaction.RollbackAsync();
throw;
}
}
}
public class OutboxMessage
{
public Guid Id { get; set; }
public string Type { get; set; }
public string Data { get; set; }
public DateTime CreatedAt { get; set; }
public DateTime? ProcessedAt { get; set; }
}
public class EventStore
{
public Guid Id { get; set; }
public Guid AggregateId { get; set; }
public int Version { get; set; }
public string EventType { get; set; }
public string EventData { get; set; }
public DateTime Timestamp { get; set; }
}
73. How do you implement database API endpoints?
Answer: Database API endpoints provide secure, controlled access to database operations through RESTful or GraphQL interfaces with proper validation, authorization, and error handling.
// Database API Controller with Repository Pattern
[ApiController]
[Route("api/[controller]")]
[Authorize]
public class UsersController : ControllerBase
{
private readonly IUserRepository _userRepository;
private readonly ILogger<UsersController> _logger;
private readonly IMapper _mapper;
public UsersController(
IUserRepository userRepository,
ILogger<UsersController> logger,
IMapper mapper)
{
_userRepository = userRepository;
_logger = logger;
_mapper = mapper;
}
[HttpGet]
public async Task<ActionResult<PagedResult<UserDto>>> GetUsers(
[FromQuery] UserQueryParameters parameters)
{
try
{
var users = await _userRepository.GetUsersAsync(parameters);
var userDtos = _mapper.Map<List<UserDto>>(users.Items);
return Ok(new PagedResult<UserDto>
{
Items = userDtos,
TotalCount = users.TotalCount,
PageNumber = users.PageNumber,
PageSize = users.PageSize
});
}
catch (Exception ex)
{
_logger.LogError(ex, "Error retrieving users");
return StatusCode(500, "Internal server error");
}
}
[HttpGet("{id:guid}")]
public async Task<ActionResult<UserDto>> GetUser(Guid id)
{
var user = await _userRepository.GetByIdAsync(id);
if (user == null)
return NotFound();
return Ok(_mapper.Map<UserDto>(user));
}
[HttpPost]
public async Task<ActionResult<UserDto>> CreateUser([FromBody] CreateUserRequest request)
{
if (!ModelState.IsValid)
return BadRequest(ModelState);
try
{
var user = _mapper.Map<User>(request);
var createdUser = await _userRepository.CreateAsync(user);
return CreatedAtAction(
nameof(GetUser),
new { id = createdUser.Id },
_mapper.Map<UserDto>(createdUser));
}
catch (ValidationException ex)
{
return BadRequest(ex.Message);
}
catch (Exception ex)
{
_logger.LogError(ex, "Error creating user");
return StatusCode(500, "Internal server error");
}
}
[HttpPut("{id:guid}")]
public async Task<ActionResult> UpdateUser(Guid id, [FromBody] UpdateUserRequest request)
{
if (!ModelState.IsValid)
return BadRequest(ModelState);
try
{
var user = await _userRepository.GetByIdAsync(id);
if (user == null)
return NotFound();
_mapper.Map(request, user);
await _userRepository.UpdateAsync(user);
return NoContent();
}
catch (Exception ex)
{
_logger.LogError(ex, "Error updating user {UserId}", id);
return StatusCode(500, "Internal server error");
}
}
[HttpDelete("{id:guid}")]
public async Task<ActionResult> DeleteUser(Guid id)
{
try
{
var success = await _userRepository.DeleteAsync(id);
if (!success)
return NotFound();
return NoContent();
}
catch (Exception ex)
{
_logger.LogError(ex, "Error deleting user {UserId}", id);
return StatusCode(500, "Internal server error");
}
}
}
// Repository Implementation
public class UserRepository : IUserRepository
{
private readonly IDbContext _dbContext;
public UserRepository(IDbContext dbContext)
{
_dbContext = dbContext;
}
public async Task<PagedResult<User>> GetUsersAsync(UserQueryParameters parameters)
{
var query = _dbContext.Users.AsQueryable();
// Apply filters
if (!string.IsNullOrEmpty(parameters.SearchTerm))
{
query = query.Where(u =>
u.FirstName.Contains(parameters.SearchTerm) ||
u.LastName.Contains(parameters.SearchTerm) ||
u.Email.Contains(parameters.SearchTerm));
}
if (parameters.IsActive.HasValue)
{
query = query.Where(u => u.IsActive == parameters.IsActive.Value);
}
// Apply sorting
query = parameters.SortBy?.ToLower() switch
{
"firstname" => parameters.SortOrder == "desc"
? query.OrderByDescending(u => u.FirstName)
: query.OrderBy(u => u.FirstName),
"lastname" => parameters.SortOrder == "desc"
? query.OrderByDescending(u => u.LastName)
: query.OrderBy(u => u.LastName),
"email" => parameters.SortOrder == "desc"
? query.OrderByDescending(u => u.Email)
: query.OrderBy(u => u.Email),
_ => query.OrderBy(u => u.CreatedAt)
};
var totalCount = await query.CountAsync();
var users = await query
.Skip((parameters.PageNumber - 1) * parameters.PageSize)
.Take(parameters.PageSize)
.ToListAsync();
return new PagedResult<User>
{
Items = users,
TotalCount = totalCount,
PageNumber = parameters.PageNumber,
PageSize = parameters.PageSize
};
}
}
74. How do you handle database data synchronization?
Answer: Database synchronization ensures data consistency across multiple databases, systems, or environments through various strategies like replication, ETL processes, or real-time synchronization.
// Database Synchronization Service
public class DatabaseSynchronizationService
{
private readonly IDbContext _sourceContext;
private readonly IDbContext _targetContext;
private readonly ILogger<DatabaseSynchronizationService> _logger;
private readonly IMessageQueue _messageQueue;
public DatabaseSynchronizationService(
IDbContext sourceContext,
IDbContext targetContext,
ILogger<DatabaseSynchronizationService> logger,
IMessageQueue messageQueue)
{
_sourceContext = sourceContext;
_targetContext = targetContext;
_logger = logger;
_messageQueue = messageQueue;
}
// Real-time synchronization using triggers and queues
public async Task SetupRealTimeSync()
{
// Create sync tracking table
await CreateSyncTrackingTable();
// Setup database triggers
await SetupDatabaseTriggers();
// Start sync processor
_ = Task.Run(ProcessSyncMessages);
}
private async Task CreateSyncTrackingTable()
{
var createTableSql = @"
CREATE TABLE IF NOT EXISTS SyncTracking (
Id UNIQUEIDENTIFIER PRIMARY KEY DEFAULT NEWID(),
TableName NVARCHAR(128) NOT NULL,
Operation NVARCHAR(10) NOT NULL,
RecordId NVARCHAR(50) NOT NULL,
SyncData NVARCHAR(MAX),
CreatedAt DATETIME2 DEFAULT GETUTCDATE(),
ProcessedAt DATETIME2 NULL,
RetryCount INT DEFAULT 0
)";
await _sourceContext.Database.ExecuteSqlRawAsync(createTableSql);
}
private async Task SetupDatabaseTriggers()
{
var triggerSql = @"
CREATE TRIGGER TR_Users_Sync
ON Users
AFTER INSERT, UPDATE, DELETE
AS
BEGIN
SET NOCOUNT ON;
DECLARE @Operation NVARCHAR(10);
DECLARE @RecordId NVARCHAR(50);
DECLARE @SyncData NVARCHAR(MAX);
IF EXISTS(SELECT * FROM inserted) AND EXISTS(SELECT * FROM deleted)
SET @Operation = 'UPDATE';
ELSE IF EXISTS(SELECT * FROM inserted)
SET @Operation = 'INSERT';
ELSE
SET @Operation = 'DELETE';
IF @Operation IN ('INSERT', 'UPDATE')
BEGIN
SELECT @RecordId = CAST(Id AS NVARCHAR(50)),
@SyncData = (SELECT * FROM inserted FOR JSON PATH, WITHOUT_ARRAY_WRAPPER)
FROM inserted;
END
ELSE
BEGIN
SELECT @RecordId = CAST(Id AS NVARCHAR(50))
FROM deleted;
END
INSERT INTO SyncTracking (TableName, Operation, RecordId, SyncData)
VALUES ('Users', @Operation, @RecordId, @SyncData);
END";
await _sourceContext.Database.ExecuteSqlRawAsync(triggerSql);
}
private async Task ProcessSyncMessages()
{
while (true)
{
try
{
var unprocessedSyncs = await _sourceContext.Set<SyncTracking>()
.Where(s => s.ProcessedAt == null && s.RetryCount < 3)
.OrderBy(s => s.CreatedAt)
.Take(50)
.ToListAsync();
foreach (var sync in unprocessedSyncs)
{
try
{
await ProcessSyncRecord(sync);
sync.ProcessedAt = DateTime.UtcNow;
}
catch (Exception ex)
{
sync.RetryCount++;
_logger.LogError(ex, "Failed to sync record {SyncId}", sync.Id);
}
}
await _sourceContext.SaveChangesAsync();
await Task.Delay(1000); // Wait 1 second before next batch
}
catch (Exception ex)
{
_logger.LogError(ex, "Error in sync processing");
await Task.Delay(5000); // Wait 5 seconds on error
}
}
}
private async Task ProcessSyncRecord(SyncTracking sync)
{
switch (sync.Operation)
{
case "INSERT":
case "UPDATE":
var userData = JsonSerializer.Deserialize<User>(sync.SyncData);
await UpsertUser(userData);
break;
case "DELETE":
await DeleteUser(Guid.Parse(sync.RecordId));
break;
}
}
private async Task UpsertUser(User user)
{
var existingUser = await _targetContext.Users
.FirstOrDefaultAsync(u => u.Id == user.Id);
if (existingUser == null)
{
_targetContext.Users.Add(user);
}
else
{
_targetContext.Entry(existingUser).CurrentValues.SetValues(user);
}
await _targetContext.SaveChangesAsync();
}
private async Task DeleteUser(Guid userId)
{
var user = await _targetContext.Users.FindAsync(userId);
if (user != null)
{
_targetContext.Users.Remove(user);
await _targetContext.SaveChangesAsync();
}
}
// Batch synchronization for large datasets
public async Task SyncBatchData(string tableName, DateTime? lastSyncTime = null)
{
var batchSize = 1000;
var offset = 0;
while (true)
{
var query = _sourceContext.Set<User>().AsQueryable();
if (lastSyncTime.HasValue)
{
query = query.Where(u => u.UpdatedAt > lastSyncTime.Value);
}
var batch = await query
.OrderBy(u => u.Id)
.Skip(offset)
.Take(batchSize)
.ToListAsync();
if (!batch.Any())
break;
await SyncBatch(batch);
offset += batchSize;
}
}
private async Task SyncBatch(List<User> users)
{
using var transaction = await _targetContext.Database.BeginTransactionAsync();
try
{
foreach (var user in users)
{
var existingUser = await _targetContext.Users
.FirstOrDefaultAsync(u => u.Id == user.Id);
if (existingUser == null)
{
_targetContext.Users.Add(user);
}
else
{
_targetContext.Entry(existingUser).CurrentValues.SetValues(user);
}
}
await _targetContext.SaveChangesAsync();
await transaction.CommitAsync();
}
catch
{
await transaction.RollbackAsync();
throw;
}
}
}
public class SyncTracking
{
public Guid Id { get; set; }
public string TableName { get; set; }
public string Operation { get; set; }
public string RecordId { get; set; }
public string SyncData { get; set; }
public DateTime CreatedAt { get; set; }
public DateTime? ProcessedAt { get; set; }
public int RetryCount { get; set; }
}
75. How do you implement database event sourcing?
Answer: Event sourcing stores all changes to application state as a sequence of events, enabling audit trails, temporal queries, and rebuilding application state from events.
// Event Sourcing Implementation
public class EventSourcingService
{
private readonly IDbContext _dbContext;
private readonly IEventStore _eventStore;
private readonly IEventPublisher _eventPublisher;
private readonly ILogger<EventSourcingService> _logger;
public EventSourcingService(
IDbContext dbContext,
IEventStore eventStore,
IEventPublisher eventPublisher,
ILogger<EventSourcingService> logger)
{
_dbContext = dbContext;
_eventStore = eventStore;
_eventPublisher = eventPublisher;
_logger = logger;
}
// Save aggregate with events
public async Task SaveAggregate<T>(T aggregate) where T : AggregateRoot
{
var events = aggregate.GetUncommittedEvents();
if (!events.Any())
return;
using var transaction = await _dbContext.Database.BeginTransactionAsync();
try
{
// Save aggregate snapshot
await SaveAggregateSnapshot(aggregate);
// Save events
foreach (var domainEvent in events)
{
await _eventStore.AppendEventAsync(domainEvent);
}
// Publish events
foreach (var domainEvent in events)
{
await _eventPublisher.PublishAsync(domainEvent);
}
aggregate.MarkEventsAsCommitted();
await transaction.CommitAsync();
}
catch
{
await transaction.RollbackAsync();
throw;
}
}
// Rebuild aggregate from events
public async Task<T> GetAggregate<T>(Guid aggregateId) where T : AggregateRoot, new()
{
// Try to get from snapshot first
var snapshot = await GetLatestSnapshot<T>(aggregateId);
var aggregate = snapshot?.Aggregate ?? new T();
// Get events after snapshot
var fromVersion = snapshot?.Version ?? 0;
var events = await _eventStore.GetEventsAsync(aggregateId, fromVersion);
// Apply events to rebuild state
foreach (var domainEvent in events)
{
aggregate.Apply(domainEvent);
}
return aggregate;
}
// Event Store Implementation
public class EventStore : IEventStore
{
private readonly IDbContext _dbContext;
public EventStore(IDbContext dbContext)
{
_dbContext = dbContext;
}
public async Task AppendEventAsync(DomainEvent domainEvent)
{
var eventRecord = new EventRecord
{
Id = Guid.NewGuid(),
AggregateId = domainEvent.AggregateId,
AggregateType = domainEvent.AggregateType,
EventType = domainEvent.GetType().Name,
EventData = JsonSerializer.Serialize(domainEvent),
Version = domainEvent.Version,
Timestamp = DateTime.UtcNow
};
_dbContext.EventRecords.Add(eventRecord);
await _dbContext.SaveChangesAsync();
}
public async Task<List<DomainEvent>> GetEventsAsync(Guid aggregateId, int fromVersion = 0)
{
var eventRecords = await _dbContext.EventRecords
.Where(e => e.AggregateId == aggregateId && e.Version > fromVersion)
.OrderBy(e => e.Version)
.ToListAsync();
var events = new List<DomainEvent>();
foreach (var record in eventRecords)
{
var eventType = Type.GetType(record.EventType);
var domainEvent = (DomainEvent)JsonSerializer.Deserialize(record.EventData, eventType);
events.Add(domainEvent);
}
return events;
}
}
// Snapshot Implementation
public class SnapshotService
{
private readonly IDbContext _dbContext;
public SnapshotService(IDbContext dbContext)
{
_dbContext = dbContext;
}
public async Task SaveAggregateSnapshot<T>(T aggregate) where T : AggregateRoot
{
var snapshot = new Snapshot
{
Id = Guid.NewGuid(),
AggregateId = aggregate.Id,
AggregateType = typeof(T).Name,
Version = aggregate.Version,
SnapshotData = JsonSerializer.Serialize(aggregate),
CreatedAt = DateTime.UtcNow
};
_dbContext.Snapshots.Add(snapshot);
await _dbContext.SaveChangesAsync();
}
public async Task<Snapshot> GetLatestSnapshot<T>(Guid aggregateId) where T : AggregateRoot
{
var snapshotRecord = await _dbContext.Snapshots
.Where(s => s.AggregateId == aggregateId && s.AggregateType == typeof(T).Name)
.OrderByDescending(s => s.Version)
.FirstOrDefaultAsync();
if (snapshotRecord == null)
return null;
var aggregate = JsonSerializer.Deserialize<T>(snapshotRecord.SnapshotData);
return new Snapshot
{
Aggregate = aggregate,
Version = snapshotRecord.Version
};
}
}
}
// Domain Models
public abstract class AggregateRoot
{
private readonly List<DomainEvent> _uncommittedEvents = new();
public Guid Id { get; protected set; }
public int Version { get; protected set; }
protected void Apply(DomainEvent domainEvent)
{
domainEvent.Version = Version + 1;
domainEvent.AggregateId = Id;
domainEvent.AggregateType = GetType().Name;
When(domainEvent);
_uncommittedEvents.Add(domainEvent);
Version++;
}
protected abstract void When(DomainEvent domainEvent);
public IEnumerable<DomainEvent> GetUncommittedEvents() => _uncommittedEvents;
public void MarkEventsAsCommitted() => _uncommittedEvents.Clear();
}
public abstract class DomainEvent
{
public Guid Id { get; set; } = Guid.NewGuid();
public Guid AggregateId { get; set; }
public string AggregateType { get; set; }
public int Version { get; set; }
public DateTime Timestamp { get; set; } = DateTime.UtcNow;
}
// Example Aggregate
public class User : AggregateRoot
{
public string FirstName { get; private set; }
public string LastName { get; private set; }
public string Email { get; private set; }
public bool IsActive { get; private set; }
public User(string firstName, string lastName, string email)
{
Id = Guid.NewGuid();
Apply(new UserCreatedEvent(firstName, lastName, email));
}
public void UpdateProfile(string firstName, string lastName)
{
Apply(new UserProfileUpdatedEvent(firstName, lastName));
}
public void Deactivate()
{
Apply(new UserDeactivatedEvent());
}
protected override void When(DomainEvent domainEvent)
{
switch (domainEvent)
{
case UserCreatedEvent e:
FirstName = e.FirstName;
LastName = e.LastName;
Email = e.Email;
IsActive = true;
break;
case UserProfileUpdatedEvent e:
FirstName = e.FirstName;
LastName = e.LastName;
break;
case UserDeactivatedEvent:
IsActive = false;
break;
}
}
}
public class UserCreatedEvent : DomainEvent
{
public string FirstName { get; set; }
public string LastName { get; set; }
public string Email { get; set; }
public UserCreatedEvent(string firstName, string lastName, string email)
{
FirstName = firstName;
LastName = lastName;
Email = email;
}
}
public class UserProfileUpdatedEvent : DomainEvent
{
public string FirstName { get; set; }
public string LastName { get; set; }
public UserProfileUpdatedEvent(string firstName, string lastName)
{
FirstName = firstName;
LastName = lastName;
}
}
public class UserDeactivatedEvent : DomainEvent { }
76. How do you handle database microservices integration?
Answer: Database microservices integration involves managing data consistency across multiple services while maintaining service boundaries and autonomy.
// Database per Service Pattern with Saga Implementation
public class OrderSagaService
{
private readonly IOrderRepository _orderRepository;
private readonly IInventoryService _inventoryService;
private readonly IPaymentService _paymentService;
private readonly IShippingService _shippingService;
private readonly ISagaCoordinator _sagaCoordinator;
public OrderSagaService(
IOrderRepository orderRepository,
IInventoryService inventoryService,
IPaymentService paymentService,
IShippingService shippingService,
ISagaCoordinator sagaCoordinator)
{
_orderRepository = orderRepository;
_inventoryService = inventoryService;
_paymentService = paymentService;
_shippingService = shippingService;
_sagaCoordinator = sagaCoordinator;
}
public async Task<OrderResult> CreateOrder(CreateOrderRequest request)
{
var sagaId = Guid.NewGuid();
var saga = new OrderSaga(sagaId, request);
try
{
// Step 1: Create order
var order = await CreateOrderStep(saga);
// Step 2: Reserve inventory
await ReserveInventoryStep(saga);
// Step 3: Process payment
await ProcessPaymentStep(saga);
// Step 4: Create shipping
await CreateShippingStep(saga);
// Mark saga as completed
await _sagaCoordinator.CompleteSagaAsync(sagaId);
return new OrderResult { Success = true, OrderId = order.Id };
}
catch (Exception ex)
{
// Compensate for completed steps
await CompensateSaga(saga);
throw;
}
}
private async Task<Order> CreateOrderStep(OrderSaga saga)
{
var order = new Order
{
Id = Guid.NewGuid(),
CustomerId = saga.Request.CustomerId,
Items = saga.Request.Items,
TotalAmount = saga.Request.Items.Sum(i => i.Price * i.Quantity),
Status = OrderStatus.Created
};
await _orderRepository.CreateAsync(order);
saga.AddStep(new SagaStep("CreateOrder", order.Id.ToString()));
return order;
}
private async Task ReserveInventoryStep(OrderSaga saga)
{
var inventoryRequest = new ReserveInventoryRequest
{
Items = saga.Request.Items
};
await _inventoryService.ReserveInventoryAsync(inventoryRequest);
saga.AddStep(new SagaStep("ReserveInventory", inventoryRequest.ToString()));
}
private async Task ProcessPaymentStep(OrderSaga saga)
{
var paymentRequest = new ProcessPaymentRequest
{
OrderId = saga.OrderId,
Amount = saga.Request.Items.Sum(i => i.Price * i.Quantity),
PaymentMethod = saga.Request.PaymentMethod
};
await _paymentService.ProcessPaymentAsync(paymentRequest);
saga.AddStep(new SagaStep("ProcessPayment", paymentRequest.ToString()));
}
private async Task CreateShippingStep(OrderSaga saga)
{
var shippingRequest = new CreateShippingRequest
{
OrderId = saga.OrderId,
ShippingAddress = saga.Request.ShippingAddress
};
await _shippingService.CreateShippingAsync(shippingRequest);
saga.AddStep(new SagaStep("CreateShipping", shippingRequest.ToString()));
}
private async Task CompensateSaga(OrderSaga saga)
{
var steps = saga.GetSteps().Reverse(); // Compensate in reverse order
foreach (var step in steps)
{
try
{
await CompensateStep(step);
}
catch (Exception ex)
{
// Log compensation failure but continue with other compensations
// In production, you might want to retry or alert
}
}
}
private async Task CompensateStep(SagaStep step)
{
switch (step.Name)
{
case "CreateShipping":
await _shippingService.CancelShippingAsync(step.Data);
break;
case "ProcessPayment":
await _paymentService.RefundPaymentAsync(step.Data);
break;
case "ReserveInventory":
await _inventoryService.ReleaseInventoryAsync(step.Data);
break;
case "CreateOrder":
await _orderRepository.CancelOrderAsync(Guid.Parse(step.Data));
break;
}
}
}
// Database Integration Service
public class DatabaseIntegrationService
{
private readonly IDbContext _dbContext;
private readonly IMessageQueue _messageQueue;
private readonly ILogger<DatabaseIntegrationService> _logger;
public DatabaseIntegrationService(
IDbContext dbContext,
IMessageQueue messageQueue,
ILogger<DatabaseIntegrationService> logger)
{
_dbContext = dbContext;
_messageQueue = messageQueue;
_logger = logger;
}
// Event-driven integration
public async Task PublishDomainEvent(DomainEvent domainEvent)
{
using var transaction = await _dbContext.Database.BeginTransactionAsync();
try
{
// Save event to local event store
var eventRecord = new EventRecord
{
Id = Guid.NewGuid(),
AggregateId = domainEvent.AggregateId,
EventType = domainEvent.GetType().Name,
EventData = JsonSerializer.Serialize(domainEvent),
Timestamp = DateTime.UtcNow
};
_dbContext.EventRecords.Add(eventRecord);
await _dbContext.SaveChangesAsync();
// Publish to message queue for other services
await _messageQueue.PublishAsync("DomainEvent", JsonSerializer.Serialize(domainEvent));
await transaction.CommitAsync();
}
catch
{
await transaction.RollbackAsync();
throw;
}
}
// Data consistency through eventual consistency
public async Task HandleIntegrationEvent(string eventType, string eventData)
{
try
{
switch (eventType)
{
case "UserCreated":
await HandleUserCreated(JsonSerializer.Deserialize<UserCreatedEvent>(eventData));
break;
case "UserUpdated":
await HandleUserUpdated(JsonSerializer.Deserialize<UserUpdatedEvent>(eventData));
break;
case "UserDeleted":
await HandleUserDeleted(JsonSerializer.Deserialize<UserDeletedEvent>(eventData));
break;
}
}
catch (Exception ex)
{
_logger.LogError(ex, "Failed to handle integration event {EventType}", eventType);
throw;
}
}
private async Task HandleUserCreated(UserCreatedEvent userEvent)
{
// Create local copy of user data for this service
var user = new User
{
Id = userEvent.UserId,
FirstName = userEvent.FirstName,
LastName = userEvent.LastName,
Email = userEvent.Email
};
_dbContext.Users.Add(user);
await _dbContext.SaveChangesAsync();
}
private async Task HandleUserUpdated(UserUpdatedEvent userEvent)
{
var user = await _dbContext.Users.FindAsync(userEvent.UserId);
if (user != null)
{
user.FirstName = userEvent.FirstName;
user.LastName = userEvent.LastName;
user.Email = userEvent.Email;
await _dbContext.SaveChangesAsync();
}
}
private async Task HandleUserDeleted(UserDeletedEvent userEvent)
{
var user = await _dbContext.Users.FindAsync(userEvent.UserId);
if (user != null)
{
_dbContext.Users.Remove(user);
await _dbContext.SaveChangesAsync();
}
}
}
// Saga Coordinator
public class SagaCoordinator : ISagaCoordinator
{
private readonly IDbContext _dbContext;
public SagaCoordinator(IDbContext dbContext)
{
_dbContext = dbContext;
}
public async Task CompleteSagaAsync(Guid sagaId)
{
var saga = await _dbContext.Sagas.FindAsync(sagaId);
if (saga != null)
{
saga.Status = SagaStatus.Completed;
saga.CompletedAt = DateTime.UtcNow;
await _dbContext.SaveChangesAsync();
}
}
public async Task FailSagaAsync(Guid sagaId, string reason)
{
var saga = await _dbContext.Sagas.FindAsync(sagaId);
if (saga != null)
{
saga.Status = SagaStatus.Failed;
saga.FailureReason = reason;
saga.FailedAt = DateTime.UtcNow;
await _dbContext.SaveChangesAsync();
}
}
}
77. How do you implement database CQRS patterns?
CQRS (Command Query Responsibility Segregation) separates read and write operations into different models and potentially different databases.
Key Concepts:
- Commands: Write operations that change state
- Queries: Read operations that retrieve data
- Event Sourcing: Store events instead of current state
- Projections: Build read models from events
Implementation Example:
// Domain Models
public class Order
{
public Guid Id { get; private set; }
public string CustomerName { get; private set; }
public decimal TotalAmount { get; private set; }
public OrderStatus Status { get; private set; }
public List<OrderItem> Items { get; private set; } = new();
public Order(Guid id, string customerName)
{
Id = id;
CustomerName = customerName;
Status = OrderStatus.Created;
}
public void AddItem(string productName, int quantity, decimal price)
{
var item = new OrderItem(productName, quantity, price);
Items.Add(item);
TotalAmount = Items.Sum(i => i.TotalPrice);
}
public void Confirm()
{
Status = OrderStatus.Confirmed;
}
}
// Commands
public abstract class Command
{
public Guid Id { get; set; }
public DateTime Timestamp { get; set; } = DateTime.UtcNow;
}
public class CreateOrderCommand : Command
{
public string CustomerName { get; set; }
public List<OrderItemDto> Items { get; set; }
}
public class ConfirmOrderCommand : Command
{
public Guid OrderId { get; set; }
}
// Command Handler
public class OrderCommandHandler : ICommandHandler<CreateOrderCommand>, ICommandHandler<ConfirmOrderCommand>
{
private readonly IEventStore _eventStore;
private readonly IEventPublisher _eventPublisher;
public OrderCommandHandler(IEventStore eventStore, IEventPublisher eventPublisher)
{
_eventStore = eventStore;
_eventPublisher = eventPublisher;
}
public async Task HandleAsync(CreateOrderCommand command)
{
var order = new Order(command.Id, command.CustomerName);
foreach (var item in command.Items)
{
order.AddItem(item.ProductName, item.Quantity, item.Price);
}
var events = order.GetUncommittedEvents();
await _eventStore.SaveEventsAsync(command.Id, events, -1);
await _eventPublisher.PublishAsync(events);
}
public async Task HandleAsync(ConfirmOrderCommand command)
{
var events = await _eventStore.GetEventsAsync(command.OrderId);
var order = Order.ReconstructFromEvents(events);
order.Confirm();
var newEvents = order.GetUncommittedEvents();
await _eventStore.SaveEventsAsync(command.OrderId, newEvents, events.Count);
await _eventPublisher.PublishAsync(newEvents);
}
}
// Queries
public class OrderQueryModel
{
public Guid Id { get; set; }
public string CustomerName { get; set; }
public decimal TotalAmount { get; set; }
public OrderStatus Status { get; set; }
public List<OrderItemQueryModel> Items { get; set; }
public DateTime CreatedDate { get; set; }
}
public class OrderItemQueryModel
{
public string ProductName { get; set; }
public int Quantity { get; set; }
public decimal Price { get; set; }
public decimal TotalPrice { get; set; }
}
// Query Handler
public class OrderQueryHandler : IQueryHandler<GetOrderQuery, OrderQueryModel>
{
private readonly IOrderReadRepository _readRepository;
public OrderQueryHandler(IOrderReadRepository readRepository)
{
_readRepository = readRepository;
}
public async Task<OrderQueryModel> HandleAsync(GetOrderQuery query)
{
return await _readRepository.GetByIdAsync(query.OrderId);
}
}
// Read Repository (Optimized for queries)
public class OrderReadRepository : IOrderReadRepository
{
private readonly IDbConnection _connection;
public OrderReadRepository(IDbConnection connection)
{
_connection = connection;
}
public async Task<OrderQueryModel> GetByIdAsync(Guid orderId)
{
const string sql = @"
SELECT o.Id, o.CustomerName, o.TotalAmount, o.Status, o.CreatedDate,
oi.ProductName, oi.Quantity, oi.Price, oi.TotalPrice
FROM Orders o
LEFT JOIN OrderItems oi ON o.Id = oi.OrderId
WHERE o.Id = @OrderId";
var orderDictionary = new Dictionary<Guid, OrderQueryModel>();
await _connection.QueryAsync<OrderQueryModel, OrderItemQueryModel, OrderQueryModel>(
sql,
(order, item) =>
{
if (!orderDictionary.TryGetValue(order.Id, out var orderEntry))
{
orderEntry = order;
orderEntry.Items = new List<OrderItemQueryModel>();
orderDictionary.Add(order.Id, orderEntry);
}
if (item != null)
orderEntry.Items.Add(item);
return orderEntry;
},
new { OrderId = orderId },
splitOn: "ProductName"
);
return orderDictionary.Values.FirstOrDefault();
}
}
// Event Store
public interface IEventStore
{
Task SaveEventsAsync(Guid aggregateId, IEnumerable<DomainEvent> events, int expectedVersion);
Task<IEnumerable<DomainEvent>> GetEventsAsync(Guid aggregateId);
}
public class SqlEventStore : IEventStore
{
private readonly IDbConnection _connection;
public SqlEventStore(IDbConnection connection)
{
_connection = connection;
}
public async Task SaveEventsAsync(Guid aggregateId, IEnumerable<DomainEvent> events, int expectedVersion)
{
using var transaction = _connection.BeginTransaction();
try
{
var currentVersion = await GetCurrentVersionAsync(aggregateId);
if (currentVersion != expectedVersion)
throw new ConcurrencyException($"Expected version {expectedVersion}, but got {currentVersion}");
foreach (var @event in events)
{
await _connection.ExecuteAsync(@"
INSERT INTO EventStore (AggregateId, Version, EventType, EventData, Timestamp)
VALUES (@AggregateId, @Version, @EventType, @EventData, @Timestamp)",
new
{
AggregateId = aggregateId,
Version = ++currentVersion,
EventType = @event.GetType().Name,
EventData = JsonSerializer.Serialize(@event),
Timestamp = DateTime.UtcNow
});
}
transaction.Commit();
}
catch
{
transaction.Rollback();
throw;
}
}
public async Task<IEnumerable<DomainEvent>> GetEventsAsync(Guid aggregateId)
{
var events = await _connection.QueryAsync(@"
SELECT EventType, EventData
FROM EventStore
WHERE AggregateId = @AggregateId
ORDER BY Version",
new { AggregateId = aggregateId });
return events.Select(e => JsonSerializer.Deserialize<DomainEvent>(e.EventData,
new JsonSerializerOptions { PropertyNameCaseInsensitive = true }));
}
private async Task<int> GetCurrentVersionAsync(Guid aggregateId)
{
var result = await _connection.QueryFirstOrDefaultAsync<int>(@"
SELECT MAX(Version) FROM EventStore WHERE AggregateId = @AggregateId",
new { AggregateId = aggregateId });
return result;
}
}
78. How do you handle database event-driven architecture?
Event-Driven Architecture uses events to trigger and communicate between decoupled services.
Key Concepts:
- Event Sourcing: Store events as the source of truth
- Event Streaming: Use message brokers for event distribution
- Event Handlers: Process events asynchronously
- Event Versioning: Handle schema evolution
Implementation Example:
// Domain Events
public abstract class DomainEvent
{
public Guid Id { get; set; } = Guid.NewGuid();
public DateTime Timestamp { get; set; } = DateTime.UtcNow;
public string EventType { get; set; }
public long Version { get; set; }
}
public class OrderCreatedEvent : DomainEvent
{
public Guid OrderId { get; set; }
public string CustomerName { get; set; }
public decimal TotalAmount { get; set; }
}
public class OrderConfirmedEvent : DomainEvent
{
public Guid OrderId { get; set; }
public DateTime ConfirmedAt { get; set; }
}
// Event Publisher
public interface IEventPublisher
{
Task PublishAsync<T>(T @event) where T : DomainEvent;
Task PublishAsync(IEnumerable<DomainEvent> events);
}
public class AzureServiceBusEventPublisher : IEventPublisher
{
private readonly ServiceBusClient _client;
private readonly ILogger<AzureServiceBusEventPublisher> _logger;
public AzureServiceBusEventPublisher(ServiceBusClient client, ILogger<AzureServiceBusEventPublisher> logger)
{
_client = client;
_logger = logger;
}
public async Task PublishAsync<T>(T @event) where T : DomainEvent
{
var sender = _client.CreateSender(GetTopicName<T>());
var message = new ServiceBusMessage(JsonSerializer.Serialize(@event))
{
MessageId = @event.Id.ToString(),
CorrelationId = @event.Id.ToString(),
ApplicationProperties =
{
["EventType"] = @event.EventType,
["Version"] = @event.Version.ToString()
}
};
await sender.SendMessageAsync(message);
_logger.LogInformation("Published event {EventType} with ID {EventId}", @event.EventType, @event.Id);
}
public async Task PublishAsync(IEnumerable<DomainEvent> events)
{
foreach (var @event in events)
{
await PublishAsync(@event);
}
}
private string GetTopicName<T>() where T : DomainEvent
{
return typeof(T).Name.ToLowerInvariant();
}
}
// Event Handlers
public interface IEventHandler<T> where T : DomainEvent
{
Task HandleAsync(T @event);
}
public class OrderCreatedEventHandler : IEventHandler<OrderCreatedEvent>
{
private readonly IEmailService _emailService;
private readonly IInventoryService _inventoryService;
private readonly ILogger<OrderCreatedEventHandler> _logger;
public OrderCreatedEventHandler(
IEmailService emailService,
IInventoryService inventoryService,
ILogger<OrderCreatedEventHandler> logger)
{
_emailService = emailService;
_inventoryService = inventoryService;
_logger = logger;
}
public async Task HandleAsync(OrderCreatedEvent @event)
{
_logger.LogInformation("Handling OrderCreatedEvent for order {OrderId}", @event.OrderId);
// Send confirmation email
await _emailService.SendOrderConfirmationAsync(@event.OrderId, @event.CustomerName);
// Update inventory
await _inventoryService.ReserveInventoryAsync(@event.OrderId);
_logger.LogInformation("Successfully handled OrderCreatedEvent for order {OrderId}", @event.OrderId);
}
}
public class OrderConfirmedEventHandler : IEventHandler<OrderConfirmedEvent>
{
private readonly IPaymentService _paymentService;
private readonly IShippingService _shippingService;
private readonly ILogger<OrderConfirmedEventHandler> _logger;
public OrderConfirmedEventHandler(
IPaymentService paymentService,
IShippingService shippingService,
ILogger<OrderConfirmedEventHandler> logger)
{
_paymentService = paymentService;
_shippingService = shippingService;
_logger = logger;
}
public async Task HandleAsync(OrderConfirmedEvent @event)
{
_logger.LogInformation("Handling OrderConfirmedEvent for order {OrderId}", @event.OrderId);
// Process payment
await _paymentService.ProcessPaymentAsync(@event.OrderId);
// Initiate shipping
await _shippingService.CreateShipmentAsync(@event.OrderId);
_logger.LogInformation("Successfully handled OrderConfirmedEvent for order {OrderId}", @event.OrderId);
}
}
// Event Consumer
public class EventConsumer : IHostedService
{
private readonly ServiceBusClient _client;
private readonly IServiceProvider _serviceProvider;
private readonly ILogger<EventConsumer> _logger;
private readonly Dictionary<string, Type> _eventTypes;
private ServiceBusProcessor _processor;
public EventConsumer(
ServiceBusClient client,
IServiceProvider serviceProvider,
ILogger<EventConsumer> logger)
{
_client = client;
_serviceProvider = serviceProvider;
_logger = logger;
_eventTypes = new Dictionary<string, Type>
{
["OrderCreatedEvent"] = typeof(OrderCreatedEvent),
["OrderConfirmedEvent"] = typeof(OrderConfirmedEvent)
};
}
public async Task StartAsync(CancellationToken cancellationToken)
{
_processor = _client.CreateProcessor("order-events", new ServiceBusProcessorOptions
{
MaxConcurrentCalls = 1,
AutoCompleteMessages = false
});
_processor.ProcessMessageAsync += ProcessMessageAsync;
_processor.ProcessErrorAsync += ProcessErrorAsync;
await _processor.StartProcessingAsync(cancellationToken);
}
private async Task ProcessMessageAsync(ProcessMessageEventArgs args)
{
try
{
var eventType = args.Message.ApplicationProperties["EventType"].ToString();
var eventData = args.Message.Body.ToString();
if (_eventTypes.TryGetValue(eventType, out var type))
{
var @event = JsonSerializer.Deserialize(eventData, type) as DomainEvent;
using var scope = _serviceProvider.CreateScope();
var handlerType = typeof(IEventHandler<>).MakeGenericType(type);
var handler = scope.ServiceProvider.GetService(handlerType);
if (handler != null)
{
await (Task)handlerType.GetMethod("HandleAsync").Invoke(handler, new[] { @event });
}
}
await args.CompleteMessageAsync(args.Message);
}
catch (Exception ex)
{
_logger.LogError(ex, "Error processing message {MessageId}", args.Message.MessageId);
await args.DeadLetterMessageAsync(args.Message);
}
}
private Task ProcessErrorAsync(ProcessErrorEventArgs args)
{
_logger.LogError(args.Exception, "Error processing message from {EntityPath}", args.EntityPath);
return Task.CompletedTask;
}
public async Task StopAsync(CancellationToken cancellationToken)
{
if (_processor != null)
{
await _processor.StopProcessingAsync(cancellationToken);
await _processor.DisposeAsync();
}
}
}
// Event Sourcing with Snapshots
public class EventSourcedAggregate<T> where T : class
{
private readonly List<DomainEvent> _uncommittedEvents = new();
private long _version;
public long Version => _version;
public IEnumerable<DomainEvent> UncommittedEvents => _uncommittedEvents.AsReadOnly();
protected void Apply(DomainEvent @event)
{
@event.Version = _version + 1;
_uncommittedEvents.Add(@event);
_version++;
}
public void MarkEventsAsCommitted()
{
_uncommittedEvents.Clear();
}
public static T ReconstructFromEvents(IEnumerable<DomainEvent> events)
{
var aggregate = Activator.CreateInstance<T>();
var eventSourced = aggregate as EventSourcedAggregate<T>;
foreach (var @event in events)
{
eventSourced.Apply(@event);
eventSourced._version = @event.Version;
}
return aggregate;
}
}
// Snapshot Support
public class Snapshot
{
public Guid AggregateId { get; set; }
public long Version { get; set; }
public string Data { get; set; }
public DateTime CreatedAt { get; set; }
}
public interface ISnapshotStore
{
Task<Snapshot> GetLatestSnapshotAsync(Guid aggregateId);
Task SaveSnapshotAsync(Snapshot snapshot);
}
public class SqlSnapshotStore : ISnapshotStore
{
private readonly IDbConnection _connection;
public SqlSnapshotStore(IDbConnection connection)
{
_connection = connection;
}
public async Task<Snapshot> GetLatestSnapshotAsync(Guid aggregateId)
{
return await _connection.QueryFirstOrDefaultAsync<Snapshot>(@"
SELECT AggregateId, Version, Data, CreatedAt
FROM Snapshots
WHERE AggregateId = @AggregateId
ORDER BY Version DESC",
new { AggregateId = aggregateId });
}
public async Task SaveSnapshotAsync(Snapshot snapshot)
{
await _connection.ExecuteAsync(@"
INSERT INTO Snapshots (AggregateId, Version, Data, CreatedAt)
VALUES (@AggregateId, @Version, @Data, @CreatedAt)",
snapshot);
}
}
79. How do you implement database API versioning?
API Versioning allows you to maintain multiple versions of your API while ensuring backward compatibility.
Key Concepts:
- URL Versioning:
/api/v1/orders,/api/v2/orders - Header Versioning:
Accept: application/vnd.company.v1+json - Query Parameter Versioning:
/api/orders?version=1 - Content Negotiation: Different response formats
Implementation Example:
// API Versioning Configuration
public class ApiVersioningConfig
{
public static void ConfigureApiVersioning(IServiceCollection services)
{
services.AddApiVersioning(options =>
{
options.DefaultApiVersion = new ApiVersion(1, 0);
options.AssumeDefaultVersionWhenUnspecified = true;
options.ReportApiVersions = true;
options.ApiVersionReader = ApiVersionReader.Combine(
new UrlSegmentApiVersionReader(),
new HeaderApiVersionReader("api-version"),
new QueryStringApiVersionReader("version")
);
});
services.AddVersionedApiExplorer(options =>
{
options.GroupNameFormat = "'v'VVV";
options.SubstituteApiVersionInUrl = true;
});
}
}
// Versioned DTOs
public class OrderDtoV1
{
public Guid Id { get; set; }
public string CustomerName { get; set; }
public decimal TotalAmount { get; set; }
public string Status { get; set; }
public DateTime CreatedDate { get; set; }
}
public class OrderDtoV2
{
public Guid Id { get; set; }
public CustomerInfoDto Customer { get; set; }
public decimal TotalAmount { get; set; }
public string Status { get; set; }
public DateTime CreatedDate { get; set; }
public List<OrderItemDtoV2> Items { get; set; }
public ShippingInfoDto Shipping { get; set; }
}
public class CustomerInfoDto
{
public string Name { get; set; }
public string Email { get; set; }
public string Phone { get; set; }
}
public class OrderItemDtoV2
{
public string ProductName { get; set; }
public int Quantity { get; set; }
public decimal UnitPrice { get; set; }
public decimal TotalPrice { get; set; }
}
public class ShippingInfoDto
{
public string Address { get; set; }
public string City { get; set; }
public string PostalCode { get; set; }
public string Country { get; set; }
}
// Versioned Controllers
[ApiController]
[Route("api/v{version:apiVersion}/[controller]")]
[ApiVersion("1.0")]
[ApiVersion("2.0")]
public class OrdersController : ControllerBase
{
private readonly IOrderService _orderService;
private readonly ILogger<OrdersController> _logger;
public OrdersController(IOrderService orderService, ILogger<OrdersController> logger)
{
_orderService = orderService;
_logger = logger;
}
[HttpGet("{id}")]
[MapToApiVersion("1.0")]
public async Task<ActionResult<OrderDtoV1>> GetOrderV1(Guid id)
{
var order = await _orderService.GetOrderAsync(id);
if (order == null)
return NotFound();
var dto = new OrderDtoV1
{
Id = order.Id,
CustomerName = order.Customer.Name,
TotalAmount = order.TotalAmount,
Status = order.Status.ToString(),
CreatedDate = order.CreatedDate
};
return Ok(dto);
}
[HttpGet("{id}")]
[MapToApiVersion("2.0")]
public async Task<ActionResult<OrderDtoV2>> GetOrderV2(Guid id)
{
var order = await _orderService.GetOrderAsync(id);
if (order == null)
return NotFound();
var dto = new OrderDtoV2
{
Id = order.Id,
Customer = new CustomerInfoDto
{
Name = order.Customer.Name,
Email = order.Customer.Email,
Phone = order.Customer.Phone
},
TotalAmount = order.TotalAmount,
Status = order.Status.ToString(),
CreatedDate = order.CreatedDate,
Items = order.Items.Select(i => new OrderItemDtoV2
{
ProductName = i.ProductName,
Quantity = i.Quantity,
UnitPrice = i.UnitPrice,
TotalPrice = i.TotalPrice
}).ToList(),
Shipping = new ShippingInfoDto
{
Address = order.ShippingAddress?.Address,
City = order.ShippingAddress?.City,
PostalCode = order.ShippingAddress?.PostalCode,
Country = order.ShippingAddress?.Country
}
};
return Ok(dto);
}
[HttpPost]
[MapToApiVersion("1.0")]
public async Task<ActionResult<OrderDtoV1>> CreateOrderV1([FromBody] CreateOrderRequestV1 request)
{
var order = await _orderService.CreateOrderAsync(new CreateOrderCommand
{
CustomerName = request.CustomerName,
Items = request.Items.Select(i => new OrderItemCommand
{
ProductName = i.ProductName,
Quantity = i.Quantity,
Price = i.Price
}).ToList()
});
var dto = new OrderDtoV1
{
Id = order.Id,
CustomerName = order.Customer.Name,
TotalAmount = order.TotalAmount,
Status = order.Status.ToString(),
CreatedDate = order.CreatedDate
};
return CreatedAtAction(nameof(GetOrderV1), new { id = order.Id, version = "1.0" }, dto);
}
[HttpPost]
[MapToApiVersion("2.0")]
public async Task<ActionResult<OrderDtoV2>> CreateOrderV2([FromBody] CreateOrderRequestV2 request)
{
var order = await _orderService.CreateOrderAsync(new CreateOrderCommand
{
Customer = new CustomerCommand
{
Name = request.Customer.Name,
Email = request.Customer.Email,
Phone = request.Customer.Phone
},
Items = request.Items.Select(i => new OrderItemCommand
{
ProductName = i.ProductName,
Quantity = i.Quantity,
Price = i.UnitPrice
}).ToList(),
ShippingAddress = request.Shipping != null ? new AddressCommand
{
Address = request.Shipping.Address,
City = request.Shipping.City,
PostalCode = request.Shipping.PostalCode,
Country = request.Shipping.Country
} : null
});
var dto = new OrderDtoV2
{
Id = order.Id,
Customer = new CustomerInfoDto
{
Name = order.Customer.Name,
Email = order.Customer.Email,
Phone = order.Customer.Phone
},
TotalAmount = order.TotalAmount,
Status = order.Status.ToString(),
CreatedDate = order.CreatedDate,
Items = order.Items.Select(i => new OrderItemDtoV2
{
ProductName = i.ProductName,
Quantity = i.Quantity,
UnitPrice = i.UnitPrice,
TotalPrice = i.TotalPrice
}).ToList(),
Shipping = order.ShippingAddress != null ? new ShippingInfoDto
{
Address = order.ShippingAddress.Address,
City = order.ShippingAddress.City,
PostalCode = order.ShippingAddress.PostalCode,
Country = order.ShippingAddress.Country
} : null
};
return CreatedAtAction(nameof(GetOrderV2), new { id = order.Id, version = "2.0" }, dto);
}
}
// Versioned Request DTOs
public class CreateOrderRequestV1
{
public string CustomerName { get; set; }
public List<OrderItemRequestV1> Items { get; set; }
}
public class OrderItemRequestV1
{
public string ProductName { get; set; }
public int Quantity { get; set; }
public decimal Price { get; set; }
}
public class CreateOrderRequestV2
{
public CustomerRequestDto Customer { get; set; }
public List<OrderItemRequestV2> Items { get; set; }
public ShippingRequestDto Shipping { get; set; }
}
public class CustomerRequestDto
{
public string Name { get; set; }
public string Email { get; set; }
public string Phone { get; set; }
}
public class OrderItemRequestV2
{
public string ProductName { get; set; }
public int Quantity { get; set; }
public decimal UnitPrice { get; set; }
}
public class ShippingRequestDto
{
public string Address { get; set; }
public string City { get; set; }
public string PostalCode { get; set; }
public string Country { get; set; }
}
// Version Migration Service
public interface IVersionMigrationService
{
Task<T> MigrateToVersionAsync<T>(object source, string targetVersion);
}
public class VersionMigrationService : IVersionMigrationService
{
private readonly ILogger<VersionMigrationService> _logger;
public VersionMigrationService(ILogger<VersionMigrationService> logger)
{
_logger = logger;
}
public async Task<T> MigrateToVersionAsync<T>(object source, string targetVersion)
{
_logger.LogInformation("Migrating object from {SourceType} to version {TargetVersion}",
source.GetType().Name, targetVersion);
// Implement migration logic based on source type and target version
return await Task.FromResult((T)source);
}
}
// API Version Documentation
[ApiController]
[Route("api/[controller]")]
public class VersionController : ControllerBase
{
[HttpGet]
public IActionResult GetVersionInfo()
{
var versionInfo = new
{
CurrentVersion = "2.0",
SupportedVersions = new[] { "1.0", "2.0" },
DeprecatedVersions = new[] { "1.0" },
MigrationGuide = new
{
FromV1ToV2 = "https://api.company.com/docs/migration-v1-to-v2",
BreakingChanges = new[]
{
"Customer information is now nested under 'customer' property",
"Order items include unit price and total price",
"Shipping information is now included in responses"
}
}
};
return Ok(versionInfo);
}
}
80. How do you handle database third-party integrations?
Third-party integrations involve connecting your database with external services, APIs, and systems.
Key Concepts:
- API Integration: REST, GraphQL, gRPC
- Message Queues: Reliable message delivery
- Data Synchronization: Keep data in sync across systems
- Error Handling: Retry logic and circuit breakers
- Security: Authentication and encryption
Implementation Example:
// Integration Configuration
public class IntegrationConfig
{
public string PaymentApiUrl { get; set; }
public string PaymentApiKey { get; set; }
public string InventoryApiUrl { get; set; }
public string InventoryApiKey { get; set; }
public string ShippingApiUrl { get; set; }
public string ShippingApiKey { get; set; }
public int RetryAttempts { get; set; } = 3;
public int RetryDelayMs { get; set; } = 1000;
}
// Base Integration Service
public abstract class BaseIntegrationService
{
protected readonly HttpClient _httpClient;
protected readonly ILogger _logger;
protected readonly IntegrationConfig _config;
protected readonly IAsyncPolicy<HttpResponseMessage> _retryPolicy;
protected BaseIntegrationService(
HttpClient httpClient,
ILogger logger,
IntegrationConfig config)
{
_httpClient = httpClient;
_logger = logger;
_config = config;
_retryPolicy = Policy<HttpResponseMessage>
.Handle<HttpRequestException>()
.Or<TimeoutException>()
.OrResult(response => !response.IsSuccessStatusCode)
.WaitAndRetryAsync(
_config.RetryAttempts,
retryAttempt => TimeSpan.FromMilliseconds(_config.RetryDelayMs * retryAttempt),
onRetry: (exception, timeSpan, retryCount, context) =>
{
_logger.LogWarning(exception.Exception,
"Retry {RetryCount} after {Delay}ms", retryCount, timeSpan.TotalMilliseconds);
}
);
}
protected async Task<T> ExecuteWithRetryAsync<T>(Func<Task<T>> operation)
{
return await _retryPolicy.ExecuteAsync(async () => await operation());
}
}
// Payment Service Integration
public interface IPaymentService
{
Task<PaymentResult> ProcessPaymentAsync(PaymentRequest request);
Task<PaymentStatus> GetPaymentStatusAsync(string paymentId);
Task<bool> RefundPaymentAsync(string paymentId, decimal amount);
}
public class PaymentService : BaseIntegrationService, IPaymentService
{
public PaymentService(
HttpClient httpClient,
ILogger<PaymentService> logger,
IntegrationConfig config) : base(httpClient, logger, config)
{
_httpClient.BaseAddress = new Uri(_config.PaymentApiUrl);
_httpClient.DefaultRequestHeaders.Add("X-API-Key", _config.PaymentApiKey);
}
public async Task<PaymentResult> ProcessPaymentAsync(PaymentRequest request)
{
return await ExecuteWithRetryAsync(async () =>
{
_logger.LogInformation("Processing payment for order {OrderId}", request.OrderId);
var response = await _httpClient.PostAsJsonAsync("/api/payments", request);
response.EnsureSuccessStatusCode();
var result = await response.Content.ReadFromJsonAsync<PaymentResult>();
_logger.LogInformation("Payment processed successfully for order {OrderId}, PaymentId: {PaymentId}",
request.OrderId, result.PaymentId);
return result;
});
}
public async Task<PaymentStatus> GetPaymentStatusAsync(string paymentId)
{
return await ExecuteWithRetryAsync(async () =>
{
var response = await _httpClient.GetAsync($"/api/payments/{paymentId}/status");
response.EnsureSuccessStatusCode();
return await response.Content.ReadFromJsonAsync<PaymentStatus>();
});
}
public async Task<bool> RefundPaymentAsync(string paymentId, decimal amount)
{
return await ExecuteWithRetryAsync(async () =>
{
var refundRequest = new RefundRequest { Amount = amount };
var response = await _httpClient.PostAsJsonAsync($"/api/payments/{paymentId}/refund", refundRequest);
response.EnsureSuccessStatusCode();
var result = await response.Content.ReadFromJsonAsync<RefundResult>();
return result.Success;
});
}
}
// Inventory Service Integration
public interface IInventoryService
{
Task<bool> ReserveInventoryAsync(Guid orderId, List<InventoryItem> items);
Task<bool> ReleaseInventoryAsync(Guid orderId);
Task<List<InventoryStatus>> CheckInventoryAsync(List<string> productIds);
}
public class InventoryService : BaseIntegrationService, IInventoryService
{
public InventoryService(
HttpClient httpClient,
ILogger<InventoryService> logger,
IntegrationConfig config) : base(httpClient, logger, config)
{
_httpClient.BaseAddress = new Uri(_config.InventoryApiUrl);
_httpClient.DefaultRequestHeaders.Add("X-API-Key", _config.InventoryApiKey);
}
public async Task<bool> ReserveInventoryAsync(Guid orderId, List<InventoryItem> items)
{
return await ExecuteWithRetryAsync(async () =>
{
_logger.LogInformation("Reserving inventory for order {OrderId}", orderId);
var request = new ReserveInventoryRequest
{
OrderId = orderId,
Items = items
};
var response = await _httpClient.PostAsJsonAsync("/api/inventory/reserve", request);
response.EnsureSuccessStatusCode();
var result = await response.Content.ReadFromJsonAsync<ReserveInventoryResult>();
_logger.LogInformation("Inventory reserved successfully for order {OrderId}", orderId);
return result.Success;
});
}
public async Task<bool> ReleaseInventoryAsync(Guid orderId)
{
return await ExecuteWithRetryAsync(async () =>
{
_logger.LogInformation("Releasing inventory for order {OrderId}", orderId);
var response = await _httpClient.PostAsync($"/api/inventory/release/{orderId}", null);
response.EnsureSuccessStatusCode();
var result = await response.Content.ReadFromJsonAsync<ReleaseInventoryResult>();
_logger.LogInformation("Inventory released successfully for order {OrderId}", orderId);
return result.Success;
});
}
public async Task<List<InventoryStatus>> CheckInventoryAsync(List<string> productIds)
{
return await ExecuteWithRetryAsync(async () =>
{
var queryString = string.Join("&", productIds.Select(id => $"productIds={id}"));
var response = await _httpClient.GetAsync($"/api/inventory/status?{queryString}");
response.EnsureSuccessStatusCode();
return await response.Content.ReadFromJsonAsync<List<InventoryStatus>>();
});
}
}
// Shipping Service Integration
public interface IShippingService
{
Task<ShipmentResult> CreateShipmentAsync(Guid orderId, ShippingRequest request);
Task<ShipmentStatus> GetShipmentStatusAsync(string trackingNumber);
Task<List<ShippingOption>> GetShippingOptionsAsync(ShippingAddress from, ShippingAddress to, List<Package> packages);
}
public class ShippingService : BaseIntegrationService, IShippingService
{
public ShippingService(
HttpClient httpClient,
ILogger<ShippingService> logger,
IntegrationConfig config) : base(httpClient, logger, config)
{
_httpClient.BaseAddress = new Uri(_config.ShippingApiUrl);
_httpClient.DefaultRequestHeaders.Add("X-API-Key", _config.ShippingApiKey);
}
public async Task<ShipmentResult> CreateShipmentAsync(Guid orderId, ShippingRequest request)
{
return await ExecuteWithRetryAsync(async () =>
{
_logger.LogInformation("Creating shipment for order {OrderId}", orderId);
var response = await _httpClient.PostAsJsonAsync("/api/shipments", request);
response.EnsureSuccessStatusCode();
var result = await response.Content.ReadFromJsonAsync<ShipmentResult>();
_logger.LogInformation("Shipment created successfully for order {OrderId}, Tracking: {TrackingNumber}",
orderId, result.TrackingNumber);
return result;
});
}
public async Task<ShipmentStatus> GetShipmentStatusAsync(string trackingNumber)
{
return await ExecuteWithRetryAsync(async () =>
{
var response = await _httpClient.GetAsync($"/api/shipments/{trackingNumber}/status");
response.EnsureSuccessStatusCode();
return await response.Content.ReadFromJsonAsync<ShipmentStatus>();
});
}
public async Task<List<ShippingOption>> GetShippingOptionsAsync(ShippingAddress from, ShippingAddress to, List<Package> packages)
{
return await ExecuteWithRetryAsync(async () =>
{
var request = new ShippingOptionsRequest
{
FromAddress = from,
ToAddress = to,
Packages = packages
};
var response = await _httpClient.PostAsJsonAsync("/api/shipping/options", request);
response.EnsureSuccessStatusCode();
return await response.Content.ReadFromJsonAsync<List<ShippingOption>>();
});
}
}
// Integration Event Handler
public class IntegrationEventHandler : IEventHandler<OrderCreatedEvent>
{
private readonly IInventoryService _inventoryService;
private readonly ILogger<IntegrationEventHandler> _logger;
public IntegrationEventHandler(IInventoryService inventoryService, ILogger<IntegrationEventHandler> logger)
{
_inventoryService = inventoryService;
_logger = logger;
}
public async Task HandleAsync(OrderCreatedEvent @event)
{
try
{
var inventoryItems = @event.Items.Select(i => new InventoryItem
{
ProductId = i.ProductId,
Quantity = i.Quantity
}).ToList();
var success = await _inventoryService.ReserveInventoryAsync(@event.OrderId, inventoryItems);
if (!success)
{
_logger.LogError("Failed to reserve inventory for order {OrderId}", @event.OrderId);
// Handle inventory reservation failure
}
}
catch (Exception ex)
{
_logger.LogError(ex, "Error handling integration event for order {OrderId}", @event.OrderId);
throw;
}
}
}
// Circuit Breaker Pattern
public class CircuitBreakerPolicy
{
private readonly IAsyncPolicy<HttpResponseMessage> _circuitBreakerPolicy;
public CircuitBreakerPolicy()
{
_circuitBreakerPolicy = Policy<HttpResponseMessage>
.Handle<HttpRequestException>()
.Or<TimeoutException>()
.CircuitBreakerAsync(
exceptionsAllowedBeforeBreaking: 3,
durationOfBreak: TimeSpan.FromSeconds(30),
onBreak: (exception, duration) =>
{
// Log circuit breaker opened
},
onReset: () =>
{
// Log circuit breaker reset
}
);
}
public IAsyncPolicy<HttpResponseMessage> GetPolicy()
{
return _circuitBreakerPolicy;
}
}
// Data Synchronization Service
public interface IDataSyncService
{
Task SyncOrderDataAsync(Guid orderId);
Task SyncCustomerDataAsync(Guid customerId);
Task SyncProductDataAsync(List<string> productIds);
}
public class DataSyncService : IDataSyncService
{
private readonly IOrderRepository _orderRepository;
private readonly ICustomerRepository _customerRepository;
private readonly IProductRepository _productRepository;
private readonly IPaymentService _paymentService;
private readonly IInventoryService _inventoryService;
private readonly IShippingService _shippingService;
private readonly ILogger<DataSyncService> _logger;
public DataSyncService(
IOrderRepository orderRepository,
ICustomerRepository customerRepository,
IProductRepository productRepository,
IPaymentService paymentService,
IInventoryService inventoryService,
IShippingService shippingService,
ILogger<DataSyncService> logger)
{
_orderRepository = orderRepository;
_customerRepository = customerRepository;
_productRepository = productRepository;
_paymentService = paymentService;
_inventoryService = inventoryService;
_shippingService = shippingService;
_logger = logger;
}
public async Task SyncOrderDataAsync(Guid orderId)
{
var order = await _orderRepository.GetByIdAsync(orderId);
if (order == null) return;
// Sync payment status
if (!string.IsNullOrEmpty(order.PaymentId))
{
var paymentStatus = await _paymentService.GetPaymentStatusAsync(order.PaymentId);
order.PaymentStatus = paymentStatus.Status;
order.PaymentDate = paymentStatus.ProcessedDate;
}
// Sync shipping status
if (!string.IsNullOrEmpty(order.TrackingNumber))
{
var shipmentStatus = await _shippingService.GetShipmentStatusAsync(order.TrackingNumber);
order.ShippingStatus = shipmentStatus.Status;
order.EstimatedDelivery = shipmentStatus.EstimatedDelivery;
}
await _orderRepository.UpdateAsync(order);
_logger.LogInformation("Order data synchronized for order {OrderId}", orderId);
}
public async Task SyncCustomerDataAsync(Guid customerId)
{
// Implement customer data synchronization logic
await Task.CompletedTask;
}
public async Task SyncProductDataAsync(List<string> productIds)
{
var inventoryStatuses = await _inventoryService.CheckInventoryAsync(productIds);
foreach (var status in inventoryStatuses)
{
var product = await _productRepository.GetByProductIdAsync(status.ProductId);
if (product != null)
{
product.StockQuantity = status.AvailableQuantity;
product.LastStockUpdate = DateTime.UtcNow;
await _productRepository.UpdateAsync(product);
}
}
_logger.LogInformation("Product data synchronized for {ProductCount} products", productIds.Count);
}
}
// Integration Health Check
public class IntegrationHealthCheck : IHealthCheck
{
private readonly IPaymentService _paymentService;
private readonly IInventoryService _inventoryService;
private readonly IShippingService _shippingService;
private readonly ILogger<IntegrationHealthCheck> _logger;
public IntegrationHealthCheck(
IPaymentService paymentService,
IInventoryService inventoryService,
IShippingService shippingService,
ILogger<IntegrationHealthCheck> logger)
{
_paymentService = paymentService;
_inventoryService = inventoryService;
_shippingService = shippingService;
_logger = logger;
}
public async Task<HealthCheckResult> CheckHealthAsync(HealthCheckContext context, CancellationToken cancellationToken = default)
{
var healthChecks = new List<Task<bool>>
{
CheckPaymentServiceHealthAsync(),
CheckInventoryServiceHealthAsync(),
CheckShippingServiceHealthAsync()
};
var results = await Task.WhenAll(healthChecks);
var allHealthy = results.All(r => r);
if (allHealthy)
{
return HealthCheckResult.Healthy("All integration services are healthy");
}
return HealthCheckResult.Unhealthy("One or more integration services are unhealthy");
}
private async Task<bool> CheckPaymentServiceHealthAsync()
{
try
{
// Perform a lightweight health check
await _paymentService.GetPaymentStatusAsync("health-check");
return true;
}
catch
{
return false;
}
}
private async Task<bool> CheckInventoryServiceHealthAsync()
{
try
{
await _inventoryService.CheckInventoryAsync(new List<string> { "health-check" });
return true;
}
catch
{
return false;
}
}
private async Task<bool> CheckShippingServiceHealthAsync()
{
try
{
await _shippingService.GetShipmentStatusAsync("health-check");
return true;
}
catch
{
return false;
}
}
}
81. How do you migrate from on-premises to cloud databases?
Answer: Database migration involves careful planning, data validation, and minimal downtime strategies.
Key Approaches: - Lift and Shift: Direct migration with minimal changes - Replatform: Optimize for cloud-native features - Refactor: Redesign for cloud-native architecture
C# Example - Migration Orchestration:
public class DatabaseMigrationOrchestrator
{
private readonly ILogger<DatabaseMigrationOrchestrator> _logger;
private readonly IDatabaseConnectionFactory _connectionFactory;
private readonly IDataValidationService _validationService;
public async Task<MigrationResult> MigrateDatabaseAsync(MigrationConfig config)
{
try
{
// Phase 1: Pre-migration validation
var validationResult = await _validationService.ValidateSourceDatabaseAsync(config.SourceConnection);
if (!validationResult.IsValid)
return MigrationResult.Failed(validationResult.Errors);
// Phase 2: Schema migration
await MigrateSchemaAsync(config);
// Phase 3: Data migration with CDC (Change Data Capture)
var migrationJob = await StartDataMigrationAsync(config);
// Phase 4: Cutover with minimal downtime
await PerformCutoverAsync(config, migrationJob);
return MigrationResult.Success();
}
catch (Exception ex)
{
_logger.LogError(ex, "Database migration failed");
await RollbackMigrationAsync(config);
return MigrationResult.Failed(new[] { ex.Message });
}
}
private async Task PerformCutoverAsync(MigrationConfig config, MigrationJob job)
{
// Stop writes to source database
await StopSourceWritesAsync(config.SourceConnection);
// Wait for replication lag to catch up
await WaitForReplicationCatchUpAsync(job);
// Switch application connections
await UpdateConnectionStringsAsync(config.TargetConnection);
// Verify data integrity
var integrityCheck = await _validationService.VerifyDataIntegrityAsync(
config.SourceConnection, config.TargetConnection);
if (!integrityCheck.IsValid)
throw new MigrationException("Data integrity check failed");
}
}
public class MigrationConfig
{
public string SourceConnection { get; set; }
public string TargetConnection { get; set; }
public MigrationStrategy Strategy { get; set; }
public TimeSpan MaxDowntime { get; set; }
public bool EnableRollback { get; set; }
}
82. How do you implement database auto-scaling?
Answer: Auto-scaling dynamically adjusts database resources based on demand patterns.
C# Example - Auto-scaling Service:
public class DatabaseAutoScalingService
{
private readonly ILogger<DatabaseAutoScalingService> _logger;
private readonly IPerformanceMonitor _performanceMonitor;
private readonly IDatabaseResourceManager _resourceManager;
private readonly IScalingPolicy _scalingPolicy;
public async Task MonitorAndScaleAsync(string databaseId)
{
var metrics = await _performanceMonitor.GetDatabaseMetricsAsync(databaseId);
var scalingDecision = await _scalingPolicy.EvaluateScalingAsync(metrics);
if (scalingDecision.ShouldScale)
{
await ExecuteScalingAsync(databaseId, scalingDecision);
}
}
private async Task ExecuteScalingAsync(string databaseId, ScalingDecision decision)
{
try
{
_logger.LogInformation($"Scaling database {databaseId} - {decision.Action}");
switch (decision.Action)
{
case ScalingAction.ScaleUp:
await _resourceManager.ScaleUpAsync(databaseId, decision.TargetTier);
break;
case ScalingAction.ScaleDown:
await _resourceManager.ScaleDownAsync(databaseId, decision.TargetTier);
break;
case ScalingAction.AddReadReplica:
await _resourceManager.AddReadReplicaAsync(databaseId);
break;
}
await _performanceMonitor.RecordScalingEventAsync(databaseId, decision);
}
catch (Exception ex)
{
_logger.LogError(ex, $"Scaling failed for database {databaseId}");
throw;
}
}
}
public class ScalingPolicy : IScalingPolicy
{
public async Task<ScalingDecision> EvaluateScalingAsync(DatabaseMetrics metrics)
{
var decision = new ScalingDecision();
// CPU-based scaling
if (metrics.CpuUtilization > 80)
{
decision.ShouldScale = true;
decision.Action = ScalingAction.ScaleUp;
decision.TargetTier = CalculateTargetTier(metrics.CurrentTier, 1);
}
else if (metrics.CpuUtilization < 20 && metrics.CurrentTier > "Basic")
{
decision.ShouldScale = true;
decision.Action = ScalingAction.ScaleDown;
decision.TargetTier = CalculateTargetTier(metrics.CurrentTier, -1);
}
// Connection-based scaling
if (metrics.ActiveConnections > metrics.MaxConnections * 0.8)
{
decision.ShouldScale = true;
decision.Action = ScalingAction.AddReadReplica;
}
return decision;
}
}
83. How do you handle database serverless architectures?
Answer: Serverless databases automatically scale compute and storage based on demand, with pay-per-use pricing.
C# Example - Serverless Database Manager:
public class ServerlessDatabaseManager
{
private readonly ILogger<ServerlessDatabaseManager> _logger;
private readonly IDatabaseClient _databaseClient;
private readonly IConnectionPool _connectionPool;
public async Task<DatabaseResponse> ExecuteQueryAsync(string databaseId, string query)
{
var connection = await _connectionPool.GetConnectionAsync(databaseId);
try
{
// Serverless databases auto-scale compute units
var result = await connection.ExecuteQueryAsync(query);
// Monitor compute consumption
await MonitorComputeConsumptionAsync(databaseId);
return result;
}
finally
{
await _connectionPool.ReturnConnectionAsync(connection);
}
}
public async Task ConfigureAutoPauseAsync(string databaseId, AutoPauseConfig config)
{
var settings = new ServerlessSettings
{
MinVcores = config.MinVcores,
MaxVcores = config.MaxVcores,
AutoPauseDelay = config.AutoPauseDelay,
MinStorageGB = config.MinStorageGB,
MaxStorageGB = config.MaxStorageGB
};
await _databaseClient.UpdateServerlessConfigurationAsync(databaseId, settings);
}
private async Task MonitorComputeConsumptionAsync(string databaseId)
{
var consumption = await _databaseClient.GetComputeConsumptionAsync(databaseId);
if (consumption.CurrentVcores > consumption.MaxVcores * 0.9)
{
_logger.LogWarning($"High compute consumption detected for {databaseId}");
}
}
}
public class AutoPauseConfig
{
public int MinVcores { get; set; } = 1;
public int MaxVcores { get; set; } = 4;
public TimeSpan AutoPauseDelay { get; set; } = TimeSpan.FromMinutes(60);
public int MinStorageGB { get; set; } = 5;
public int MaxStorageGB { get; set; } = 100;
}
84. How do you implement database multi-region deployments?
Answer: Multi-region deployments provide high availability, disaster recovery, and reduced latency.
C# Example - Multi-Region Database Manager:
public class MultiRegionDatabaseManager
{
private readonly ILogger<MultiRegionDatabaseManager> _logger;
private readonly IRegionManager _regionManager;
private readonly ILoadBalancer _loadBalancer;
private readonly IReplicationManager _replicationManager;
public async Task<MultiRegionDeployment> DeployMultiRegionAsync(MultiRegionConfig config)
{
var deployment = new MultiRegionDeployment();
// Deploy primary region
deployment.PrimaryRegion = await DeployRegionAsync(config.PrimaryRegion, true);
// Deploy secondary regions
foreach (var region in config.SecondaryRegions)
{
var secondaryRegion = await DeployRegionAsync(region, false);
deployment.SecondaryRegions.Add(secondaryRegion);
}
// Configure replication
await ConfigureReplicationAsync(deployment, config.ReplicationConfig);
// Configure global load balancing
await ConfigureGlobalLoadBalancingAsync(deployment, config.LoadBalancingConfig);
return deployment;
}
private async Task<RegionDeployment> DeployRegionAsync(RegionConfig config, bool isPrimary)
{
var deployment = new RegionDeployment
{
Region = config.Region,
IsPrimary = isPrimary
};
// Deploy database instance
deployment.DatabaseId = await _regionManager.DeployDatabaseAsync(config);
// Configure read replicas if needed
if (config.ReadReplicas > 0)
{
deployment.ReadReplicas = await DeployReadReplicasAsync(
deployment.DatabaseId, config.ReadReplicas);
}
return deployment;
}
public async Task<DatabaseResponse> ExecuteQueryAsync(string query, QueryOptions options = null)
{
var region = await _loadBalancer.SelectOptimalRegionAsync(options?.UserLocation);
var connection = await GetConnectionForRegionAsync(region);
try
{
var result = await connection.ExecuteQueryAsync(query);
// Check if we need to route to primary for writes
if (IsWriteOperation(query) && !region.IsPrimary)
{
var primaryResult = await ExecuteOnPrimaryAsync(query);
return primaryResult;
}
return result;
}
catch (Exception ex)
{
_logger.LogError(ex, $"Query failed in region {region.Region}");
// Failover to another region if available
return await FailoverAndRetryAsync(query, region);
}
}
private async Task<DatabaseResponse> FailoverAndRetryAsync(string query, RegionDeployment failedRegion)
{
var availableRegions = await GetAvailableRegionsAsync();
var backupRegion = availableRegions.FirstOrDefault(r => r.Region != failedRegion.Region);
if (backupRegion != null)
{
_logger.LogInformation($"Failing over to region {backupRegion.Region}");
return await ExecuteQueryAsync(query, new QueryOptions { PreferredRegion = backupRegion.Region });
}
throw new NoAvailableRegionException("No available regions for failover");
}
}
85. How do you handle database containerization?
Answer: Containerization provides consistency, portability, and easier deployment management.
C# Example - Containerized Database Manager:
public class ContainerizedDatabaseManager
{
private readonly ILogger<ContainerizedDatabaseManager> _logger;
private readonly IDockerClient _dockerClient;
private readonly IVolumeManager _volumeManager;
private readonly INetworkManager _networkManager;
public async Task<ContainerizedDatabase> DeployDatabaseAsync(DatabaseContainerConfig config)
{
var database = new ContainerizedDatabase();
try
{
// Create persistent volume
var volume = await _volumeManager.CreateVolumeAsync(config.VolumeName, config.StorageSize);
// Create network
var network = await _networkManager.CreateNetworkAsync(config.NetworkName);
// Deploy database container
var containerConfig = new ContainerConfig
{
Image = config.DatabaseImage,
Name = config.ContainerName,
Environment = config.EnvironmentVariables,
PortBindings = new Dictionary<string, string>
{
{ config.DatabasePort, config.HostPort }
},
VolumeBindings = new Dictionary<string, string>
{
{ volume.Name, config.DataPath }
},
NetworkMode = network.Name,
RestartPolicy = RestartPolicy.Always
};
database.ContainerId = await _dockerClient.CreateContainerAsync(containerConfig);
await _dockerClient.StartContainerAsync(database.ContainerId);
// Wait for database to be ready
await WaitForDatabaseReadyAsync(database.ContainerId, config.HealthCheckEndpoint);
// Initialize database
await InitializeDatabaseAsync(database.ContainerId, config.InitializationScripts);
database.Status = DatabaseStatus.Running;
return database;
}
catch (Exception ex)
{
_logger.LogError(ex, "Failed to deploy containerized database");
await CleanupFailedDeploymentAsync(database);
throw;
}
}
public async Task<DatabaseResponse> ExecuteQueryAsync(string containerId, string query)
{
var execConfig = new ExecConfig
{
Cmd = new[] { "sqlcmd", "-Q", query },
AttachStdout = true,
AttachStderr = true
};
var execResult = await _dockerClient.ExecAsync(containerId, execConfig);
if (execResult.ExitCode != 0)
{
throw new DatabaseExecutionException($"Query failed: {execResult.Stderr}");
}
return new DatabaseResponse { Data = execResult.Stdout };
}
public async Task BackupDatabaseAsync(string containerId, string backupPath)
{
var backupCommand = $"BACKUP DATABASE [{GetDatabaseName()}] TO DISK = '{backupPath}'";
await ExecuteQueryAsync(containerId, backupCommand);
// Copy backup file from container to host
await _dockerClient.CopyFromContainerAsync(containerId, backupPath, backupPath);
}
public async Task RestoreDatabaseAsync(string containerId, string backupPath)
{
// Copy backup file to container
await _dockerClient.CopyToContainerAsync(containerId, backupPath, backupPath);
var restoreCommand = $"RESTORE DATABASE [{GetDatabaseName()}] FROM DISK = '{backupPath}'";
await ExecuteQueryAsync(containerId, restoreCommand);
}
}
public class DatabaseContainerConfig
{
public string DatabaseImage { get; set; } = "mcr.microsoft.com/mssql/server:2019-latest";
public string ContainerName { get; set; }
public string VolumeName { get; set; }
public string NetworkName { get; set; }
public string DatabasePort { get; set; } = "1433";
public string HostPort { get; set; }
public string DataPath { get; set; } = "/var/opt/mssql";
public Dictionary<string, string> EnvironmentVariables { get; set; }
public string HealthCheckEndpoint { get; set; }
public List<string> InitializationScripts { get; set; }
public int StorageSize { get; set; } = 10; // GB
}
86. How do you implement database infrastructure as code?
Answer: Infrastructure as Code (IaC) manages database infrastructure through declarative configuration files.
C# Example - Infrastructure as Code Manager:
public class DatabaseInfrastructureManager
{
private readonly ILogger<DatabaseInfrastructureManager> _logger;
private readonly ITerraformClient _terraformClient;
private readonly IARMClient _armClient;
private readonly IConfigurationValidator _configValidator;
public async Task<InfrastructureDeployment> DeployInfrastructureAsync(InfrastructureConfig config)
{
var deployment = new InfrastructureDeployment();
try
{
// Validate configuration
var validationResult = await _configValidator.ValidateAsync(config);
if (!validationResult.IsValid)
{
throw new ConfigurationValidationException(validationResult.Errors);
}
// Generate infrastructure code
var infrastructureCode = await GenerateInfrastructureCodeAsync(config);
// Deploy using Terraform
if (config.UseTerraform)
{
deployment = await DeployWithTerraformAsync(infrastructureCode, config);
}
// Deploy using ARM templates
else
{
deployment = await DeployWithARMAsync(infrastructureCode, config);
}
// Configure monitoring and alerting
await ConfigureMonitoringAsync(deployment, config.MonitoringConfig);
return deployment;
}
catch (Exception ex)
{
_logger.LogError(ex, "Infrastructure deployment failed");
await RollbackInfrastructureAsync(deployment);
throw;
}
}
private async Task<string> GenerateInfrastructureCodeAsync(InfrastructureConfig config)
{
var template = new InfrastructureTemplate();
// Generate Terraform configuration
if (config.UseTerraform)
{
return await GenerateTerraformConfigAsync(config);
}
// Generate ARM template
else
{
return await GenerateARMTemplateAsync(config);
}
}
private async Task<string> GenerateTerraformConfigAsync(InfrastructureConfig config)
{
var terraformConfig = new StringBuilder();
// Provider configuration
terraformConfig.AppendLine(@"
terraform {
required_providers {
azurerm = {
source = ""hashicorp/azurerm""
version = ""~>3.0""
}
}
}
provider ""azurerm"" {
features {}
}");
// Resource group
terraformConfig.AppendLine($@"
resource ""azurerm_resource_group"" ""rg"" {{
name = ""{config.ResourceGroupName}""
location = ""{config.Location}""
}}");
// SQL Server
terraformConfig.AppendLine($@"
resource ""azurerm_sql_server"" ""server"" {{
name = ""{config.ServerName}""
resource_group_name = azurerm_resource_group.rg.name
location = azurerm_resource_group.rg.location
version = ""12.0""
administrator_login = ""{config.AdminUsername}""
administrator_login_password = ""{config.AdminPassword}""
}}");
// Database
terraformConfig.AppendLine($@"
resource ""azurerm_sql_database"" ""database"" {{
name = ""{config.DatabaseName}""
resource_group_name = azurerm_resource_group.rg.name
server_name = azurerm_sql_server.server.name
location = azurerm_resource_group.rg.location
sku_name = ""{config.SkuName}""
tags = {{
Environment = ""{config.Environment}""
Project = ""{config.ProjectName}""
}}
}}");
return terraformConfig.ToString();
}
private async Task<InfrastructureDeployment> DeployWithTerraformAsync(string config, InfrastructureConfig infraConfig)
{
// Write Terraform configuration to file
var configPath = Path.Combine(Path.GetTempPath(), "terraform.tf");
await File.WriteAllTextAsync(configPath, config);
// Initialize Terraform
await _terraformClient.InitAsync(configPath);
// Plan deployment
var plan = await _terraformClient.PlanAsync(configPath);
// Apply deployment
var result = await _terraformClient.ApplyAsync(configPath);
return new InfrastructureDeployment
{
DeploymentId = result.DeploymentId,
Resources = result.Resources,
Status = InfrastructureStatus.Deployed
};
}
public async Task<InfrastructureState> GetInfrastructureStateAsync(string deploymentId)
{
return await _terraformClient.ShowAsync(deploymentId);
}
public async Task DestroyInfrastructureAsync(string deploymentId)
{
await _terraformClient.DestroyAsync(deploymentId);
}
}
public class InfrastructureConfig
{
public bool UseTerraform { get; set; } = true;
public string ResourceGroupName { get; set; }
public string Location { get; set; }
public string ServerName { get; set; }
public string DatabaseName { get; set; }
public string AdminUsername { get; set; }
public string AdminPassword { get; set; }
public string SkuName { get; set; } = "Basic";
public string Environment { get; set; }
public string ProjectName { get; set; }
public MonitoringConfig MonitoringConfig { get; set; }
}
87. How do you handle database cloud-native features?
Answer: Cloud-native features include managed services, auto-scaling, global distribution, and serverless capabilities.
C# Example - Cloud-Native Database Service:
public class CloudNativeDatabaseService
{
private readonly ILogger<CloudNativeDatabaseService> _logger;
private readonly ICloudDatabaseClient _databaseClient;
private readonly IGlobalDistributionManager _distributionManager;
private readonly IAutoScalingManager _scalingManager;
public async Task<CloudDatabase> CreateCloudNativeDatabaseAsync(CloudDatabaseConfig config)
{
var database = new CloudDatabase();
try
{
// Create managed database instance
database.Id = await _databaseClient.CreateDatabaseAsync(new DatabaseCreateRequest
{
Name = config.Name,
Sku = config.Sku,
Location = config.PrimaryLocation,
BackupRetentionDays = config.BackupRetentionDays,
GeoRedundantBackup = config.GeoRedundantBackup,
AutoScaling = config.AutoScaling
});
// Configure global distribution if needed
if (config.GlobalDistribution?.Enabled == true)
{
await ConfigureGlobalDistributionAsync(database.Id, config.GlobalDistribution);
}
// Configure auto-scaling
if (config.AutoScaling?.Enabled == true)
{
await ConfigureAutoScalingAsync(database.Id, config.AutoScaling);
}
// Configure advanced features
await ConfigureAdvancedFeaturesAsync(database.Id, config.AdvancedFeatures);
return database;
}
catch (Exception ex)
{
_logger.LogError(ex, "Failed to create cloud-native database");
await CleanupFailedDatabaseAsync(database.Id);
throw;
}
}
private async Task ConfigureGlobalDistributionAsync(string databaseId, GlobalDistributionConfig config)
{
foreach (var region in config.Regions)
{
await _distributionManager.AddRegionAsync(databaseId, new RegionConfig
{
Region = region,
ReplicationType = config.ReplicationType,
ConsistencyLevel = config.ConsistencyLevel
});
}
}
private async Task ConfigureAutoScalingAsync(string databaseId, AutoScalingConfig config)
{
await _scalingManager.ConfigureAutoScalingAsync(databaseId, new AutoScalingSettings
{
MinCapacity = config.MinCapacity,
MaxCapacity = config.MaxCapacity,
ScaleUpThreshold = config.ScaleUpThreshold,
ScaleDownThreshold = config.ScaleDownThreshold,
ScaleUpCooldown = config.ScaleUpCooldown,
ScaleDownCooldown = config.ScaleDownCooldown
});
}
private async Task ConfigureAdvancedFeaturesAsync(string databaseId, AdvancedFeaturesConfig config)
{
// Configure advanced security features
if (config.AdvancedThreatProtection)
{
await _databaseClient.EnableAdvancedThreatProtectionAsync(databaseId);
}
// Configure data masking
if (config.DataMasking?.Enabled == true)
{
await ConfigureDataMaskingAsync(databaseId, config.DataMasking);
}
// Configure audit logging
if (config.AuditLogging?.Enabled == true)
{
await ConfigureAuditLoggingAsync(databaseId, config.AuditLogging);
}
// Configure encryption
if (config.Encryption?.Enabled == true)
{
await ConfigureEncryptionAsync(databaseId, config.Encryption);
}
}
public async Task<QueryResult> ExecuteQueryAsync(string databaseId, string query, QueryOptions options = null)
{
// Use nearest region for read operations
var region = await _distributionManager.GetNearestRegionAsync(options?.UserLocation);
var result = await _databaseClient.ExecuteQueryAsync(databaseId, query, new QueryRequest
{
Region = region,
ConsistencyLevel = options?.ConsistencyLevel ?? ConsistencyLevel.Session,
MaxRetries = options?.MaxRetries ?? 3
});
return result;
}
public async Task<BackupResult> CreateBackupAsync(string databaseId, BackupConfig config)
{
return await _databaseClient.CreateBackupAsync(databaseId, new BackupRequest
{
BackupType = config.BackupType,
RetentionDays = config.RetentionDays,
GeoRedundant = config.GeoRedundant,
Compression = config.Compression
});
}
}
public class CloudDatabaseConfig
{
public string Name { get; set; }
public string Sku { get; set; }
public string PrimaryLocation { get; set; }
public int BackupRetentionDays { get; set; } = 7;
public bool GeoRedundantBackup { get; set; } = true;
public AutoScalingConfig AutoScaling { get; set; }
public GlobalDistributionConfig GlobalDistribution { get; set; }
public AdvancedFeaturesConfig AdvancedFeatures { get; set; }
}
public class GlobalDistributionConfig
{
public bool Enabled { get; set; }
public List<string> Regions { get; set; }
public ReplicationType ReplicationType { get; set; } = ReplicationType.Synchronous;
public ConsistencyLevel ConsistencyLevel { get; set; } = ConsistencyLevel.Session;
}
88. How do you implement database hybrid cloud strategies?
Answer: Hybrid cloud combines on-premises and cloud databases for optimal performance, cost, and compliance.
C# Example - Hybrid Cloud Database Manager:
public class HybridCloudDatabaseManager
{
private readonly ILogger<HybridCloudDatabaseManager> _logger;
private readonly IOnPremisesDatabaseClient _onPremisesClient;
private readonly ICloudDatabaseClient _cloudClient;
private readonly IDataSyncManager _syncManager;
private readonly ILoadBalancer _loadBalancer;
public async Task<HybridDatabaseSetup> SetupHybridDatabaseAsync(HybridConfig config)
{
var setup = new HybridDatabaseSetup();
try
{
// Configure on-premises database
setup.OnPremisesDatabase = await ConfigureOnPremisesDatabaseAsync(config.OnPremisesConfig);
// Configure cloud database
setup.CloudDatabase = await ConfigureCloudDatabaseAsync(config.CloudConfig);
// Setup data synchronization
await SetupDataSynchronizationAsync(setup, config.SyncConfig);
// Configure load balancing
await ConfigureLoadBalancingAsync(setup, config.LoadBalancingConfig);
// Setup disaster recovery
await SetupDisasterRecoveryAsync(setup, config.DRConfig);
return setup;
}
catch (Exception ex)
{
_logger.LogError(ex, "Hybrid database setup failed");
await CleanupFailedSetupAsync(setup);
throw;
}
}
private async Task<OnPremisesDatabase> ConfigureOnPremisesDatabaseAsync(OnPremisesConfig config)
{
var database = new OnPremisesDatabase();
// Configure high availability
if (config.HighAvailability?.Enabled == true)
{
await ConfigureOnPremisesHAAsync(config.HighAvailability);
}
// Configure backup
await ConfigureOnPremisesBackupAsync(config.BackupConfig);
// Configure monitoring
await ConfigureOnPremisesMonitoringAsync(config.MonitoringConfig);
return database;
}
private async Task<CloudDatabase> ConfigureCloudDatabaseAsync(CloudConfig config)
{
var database = new CloudDatabase();
// Create cloud database
database.Id = await _cloudClient.CreateDatabaseAsync(new DatabaseCreateRequest
{
Name = config.Name,
Sku = config.Sku,
Location = config.Location,
BackupRetentionDays = config.BackupRetentionDays
});
// Configure geo-replication if needed
if (config.GeoReplication?.Enabled == true)
{
await ConfigureGeoReplicationAsync(database.Id, config.GeoReplication);
}
return database;
}
private async Task SetupDataSynchronizationAsync(HybridDatabaseSetup setup, SyncConfig config)
{
// Configure data sync based on strategy
switch (config.Strategy)
{
case SyncStrategy.RealTime:
await SetupRealTimeSyncAsync(setup, config);
break;
case SyncStrategy.Batch:
await SetupBatchSyncAsync(setup, config);
break;
case SyncStrategy.EventDriven:
await SetupEventDrivenSyncAsync(setup, config);
break;
}
}
private async Task SetupRealTimeSyncAsync(HybridDatabaseSetup setup, SyncConfig config)
{
// Setup Change Data Capture (CDC)
await _syncManager.EnableCDCAsync(setup.OnPremisesDatabase.Id);
// Configure real-time replication
await _syncManager.ConfigureRealTimeReplicationAsync(new ReplicationConfig
{
SourceDatabase = setup.OnPremisesDatabase.Id,
TargetDatabase = setup.CloudDatabase.Id,
Tables = config.Tables,
ConflictResolution = config.ConflictResolution
});
}
public async Task<QueryResult> ExecuteQueryAsync(string query, QueryOptions options = null)
{
// Determine optimal database based on query type and data location
var targetDatabase = await DetermineOptimalDatabaseAsync(query, options);
try
{
if (targetDatabase.IsOnPremises)
{
return await _onPremisesClient.ExecuteQueryAsync(targetDatabase.Id, query);
}
else
{
return await _cloudClient.ExecuteQueryAsync(targetDatabase.Id, query);
}
}
catch (Exception ex)
{
_logger.LogError(ex, $"Query failed on {targetDatabase.Type}");
// Failover to other database
return await FailoverAndRetryAsync(query, targetDatabase);
}
}
private async Task<DatabaseTarget> DetermineOptimalDatabaseAsync(string query, QueryOptions options)
{
var analysis = await AnalyzeQueryAsync(query);
// Route to on-premises for sensitive data
if (analysis.ContainsSensitiveData)
{
return new DatabaseTarget { IsOnPremises = true, Id = GetOnPremisesDatabaseId() };
}
// Route to cloud for heavy analytics
if (analysis.IsAnalyticsQuery)
{
return new DatabaseTarget { IsOnPremises = false, Id = GetCloudDatabaseId() };
}
// Use load balancer for other queries
return await _loadBalancer.SelectOptimalDatabaseAsync(analysis, options?.UserLocation);
}
public async Task<MigrationResult> MigrateDataAsync(MigrationConfig config)
{
var migration = new DataMigration();
try
{
// Phase 1: Initial data load
await migration.LoadInitialDataAsync(config.SourceDatabase, config.TargetDatabase);
// Phase 2: Setup change tracking
await migration.SetupChangeTrackingAsync(config.SourceDatabase);
// Phase 3: Sync changes
await migration.SyncChangesAsync(config.SourceDatabase, config.TargetDatabase);
// Phase 4: Verify data integrity
var integrityCheck = await migration.VerifyDataIntegrityAsync(
config.SourceDatabase, config.TargetDatabase);
if (!integrityCheck.IsValid)
{
throw new DataIntegrityException("Data integrity check failed");
}
return MigrationResult.Success();
}
catch (Exception ex)
{
_logger.LogError(ex, "Data migration failed");
await migration.RollbackAsync();
return MigrationResult.Failed(new[] { ex.Message });
}
}
}
public class HybridConfig
{
public OnPremisesConfig OnPremisesConfig { get; set; }
public CloudConfig CloudConfig { get; set; }
public SyncConfig SyncConfig { get; set; }
public LoadBalancingConfig LoadBalancingConfig { get; set; }
public DisasterRecoveryConfig DRConfig { get; set; }
}
public class SyncConfig
{
public SyncStrategy Strategy { get; set; }
public List<string> Tables { get; set; }
public ConflictResolutionStrategy ConflictResolution { get; set; }
public TimeSpan SyncInterval { get; set; }
public bool EnableCompression { get; set; } = true;
}
89. How do you handle database cloud cost optimization?
Answer: Cloud cost optimization involves right-sizing, reserved capacity, and monitoring usage patterns.
C# Example - Cloud Cost Optimization Service:
public class CloudCostOptimizationService
{
private readonly ILogger<CloudCostOptimizationService> _logger;
private readonly ICostAnalyzer _costAnalyzer;
private readonly IResourceOptimizer _resourceOptimizer;
private readonly IReservedCapacityManager _reservedCapacityManager;
public async Task<CostOptimizationReport> AnalyzeAndOptimizeCostsAsync(string subscriptionId)
{
var report = new CostOptimizationReport();
try
{
// Analyze current costs
var costAnalysis = await _costAnalyzer.AnalyzeCostsAsync(subscriptionId);
report.CurrentCosts = costAnalysis;
// Identify optimization opportunities
var opportunities = await IdentifyOptimizationOpportunitiesAsync(costAnalysis);
report.OptimizationOpportunities = opportunities;
// Generate recommendations
var recommendations = await GenerateRecommendationsAsync(opportunities);
report.Recommendations = recommendations;
// Calculate potential savings
report.PotentialSavings = await CalculatePotentialSavingsAsync(recommendations);
return report;
}
catch (Exception ex)
{
_logger.LogError(ex, "Cost optimization analysis failed");
throw;
}
}
private async Task<List<OptimizationOpportunity>> IdentifyOptimizationOpportunitiesAsync(CostAnalysis analysis)
{
var opportunities = new List<OptimizationOpportunity>();
// Check for underutilized resources
foreach (var resource in analysis.Resources)
{
if (resource.Utilization < 20)
{
opportunities.Add(new OptimizationOpportunity
{
Type = OptimizationType.RightSizing,
ResourceId = resource.Id,
CurrentSku = resource.Sku,
RecommendedSku = await _resourceOptimizer.GetRecommendedSkuAsync(resource),
PotentialSavings = await CalculateRightSizingSavingsAsync(resource)
});
}
}
// Check for reserved capacity opportunities
var reservedCapacityOpportunities = await _reservedCapacityManager.GetReservedCapacityOpportunitiesAsync(
analysis.Resources);
opportunities.AddRange(reservedCapacityOpportunities);
// Check for storage optimization
var storageOpportunities = await AnalyzeStorageOptimizationAsync(analysis.StorageCosts);
opportunities.AddRange(storageOpportunities);
return opportunities;
}
private async Task<List<CostRecommendation>> GenerateRecommendationsAsync(List<OptimizationOpportunity> opportunities)
{
var recommendations = new List<CostRecommendation>();
foreach (var opportunity in opportunities)
{
switch (opportunity.Type)
{
case OptimizationType.RightSizing:
recommendations.Add(await GenerateRightSizingRecommendationAsync(opportunity));
break;
case OptimizationType.ReservedCapacity:
recommendations.Add(await GenerateReservedCapacityRecommendationAsync(opportunity));
break;
case OptimizationType.StorageOptimization:
recommendations.Add(await GenerateStorageOptimizationRecommendationAsync(opportunity));
break;
case OptimizationType.Serverless:
recommendations.Add(await GenerateServerlessRecommendationAsync(opportunity));
break;
}
}
return recommendations.OrderByDescending(r => r.PotentialSavings).ToList();
}
private async Task<CostRecommendation> GenerateRightSizingRecommendationAsync(OptimizationOpportunity opportunity)
{
return new CostRecommendation
{
Type = RecommendationType.RightSizing,
ResourceId = opportunity.ResourceId,
Action = $"Downgrade from {opportunity.CurrentSku} to {opportunity.RecommendedSku}",
PotentialSavings = opportunity.PotentialSavings,
RiskLevel = RiskLevel.Low,
ImplementationSteps = new List<string>
{
"Monitor performance for 1 week",
"Create backup",
"Scale down during maintenance window",
"Verify performance after scaling"
}
};
}
private async Task<CostRecommendation> GenerateReservedCapacityRecommendationAsync(OptimizationOpportunity opportunity)
{
var reservedCapacityOptions = await _reservedCapacityManager.GetReservedCapacityOptionsAsync(
opportunity.ResourceId, opportunity.CurrentSku);
return new CostRecommendation
{
Type = RecommendationType.ReservedCapacity,
ResourceId = opportunity.ResourceId,
Action = $"Purchase {reservedCapacityOptions.Term} year reserved capacity",
PotentialSavings = reservedCapacityOptions.Savings,
RiskLevel = RiskLevel.Low,
ImplementationSteps = new List<string>
{
"Analyze usage patterns",
"Select appropriate term",
"Purchase reserved capacity",
"Apply to resources"
}
};
}
public async Task<OptimizationResult> ApplyOptimizationAsync(CostRecommendation recommendation)
{
try
{
switch (recommendation.Type)
{
case RecommendationType.RightSizing:
return await ApplyRightSizingAsync(recommendation);
case RecommendationType.ReservedCapacity:
return await ApplyReservedCapacityAsync(recommendation);
case RecommendationType.StorageOptimization:
return await ApplyStorageOptimizationAsync(recommendation);
case RecommendationType.Serverless:
return await ApplyServerlessMigrationAsync(recommendation);
default:
throw new NotSupportedException($"Recommendation type {recommendation.Type} not supported");
}
}
catch (Exception ex)
{
_logger.LogError(ex, $"Failed to apply optimization: {recommendation.Type}");
return OptimizationResult.Failed(ex.Message);
}
}
private async Task<OptimizationResult> ApplyRightSizingAsync(CostRecommendation recommendation)
{
// Validate recommendation
var validation = await ValidateRightSizingRecommendationAsync(recommendation);
if (!validation.IsValid)
{
return OptimizationResult.Failed(validation.Errors);
}
// Create backup
await CreateBackupAsync(recommendation.ResourceId);
// Scale down resource
await _resourceOptimizer.ScaleResourceAsync(recommendation.ResourceId, recommendation.NewSku);
// Verify performance
var performanceCheck = await VerifyPerformanceAfterScalingAsync(recommendation.ResourceId);
if (!performanceCheck.IsAcceptable)
{
await RollbackScalingAsync(recommendation.ResourceId, recommendation.OriginalSku);
return OptimizationResult.Failed("Performance degradation detected after scaling");
}
return OptimizationResult.Success();
}
public async Task<CostMonitoringDashboard> CreateCostMonitoringDashboardAsync(string subscriptionId)
{
var dashboard = new CostMonitoringDashboard();
// Set up cost alerts
await SetupCostAlertsAsync(subscriptionId);
// Create cost tracking metrics
dashboard.Metrics = await CreateCostMetricsAsync(subscriptionId);
// Set up automated optimization
await SetupAutomatedOptimizationAsync(subscriptionId);
return dashboard;
}
private async Task SetupCostAlertsAsync(string subscriptionId)
{
var alertRules = new List<CostAlertRule>
{
new CostAlertRule
{
Name = "High Daily Cost",
Threshold = 100, // USD
Period = TimeSpan.FromDays(1),
Action = AlertAction.Email
},
new CostAlertRule
{
Name = "Unusual Cost Spike",
Threshold = 50, // % increase
Period = TimeSpan.FromHours(1),
Action = AlertAction.SMS
}
};
foreach (var rule in alertRules)
{
await _costAnalyzer.CreateCostAlertAsync(subscriptionId, rule);
}
}
}
public class CostAnalysis
{
public decimal TotalCost { get; set; }
public List<ResourceCost> Resources { get; set; }
public List<StorageCost> StorageCosts { get; set; }
public List<NetworkCost> NetworkCosts { get; set; }
public CostTrend Trend { get; set; }
}
public class OptimizationOpportunity
{
public OptimizationType Type { get; set; }
public string ResourceId { get; set; }
public string CurrentSku { get; set; }
public string RecommendedSku { get; set; }
public decimal PotentialSavings { get; set; }
public RiskLevel RiskLevel { get; set; }
}
90. Database Cloud Security Best Practices
Implementation with SQL Examples:
-- 1. Encrypt sensitive data at rest
CREATE TABLE users (
user_id INT PRIMARY KEY,
username VARCHAR(50),
email VARCHAR(100),
password_hash VARBINARY(256) ENCRYPTED, -- Column-level encryption
ssn VARCHAR(11) ENCRYPTED,
created_date DATETIME2 DEFAULT GETDATE()
);
-- 2. Implement Row-Level Security (RLS)
CREATE TABLE employee_data (
employee_id INT,
department_id INT,
salary DECIMAL(10,2),
manager_id INT
);
-- Enable RLS
ALTER TABLE employee_data ENABLE ROW LEVEL SECURITY;
-- Create security policy
CREATE SECURITY POLICY EmployeeDataPolicy
ON employee_data
FOR ALL
USING (
department_id IN (
SELECT department_id
FROM user_departments
WHERE user_id = SYSTEM_USER_ID()
)
OR
manager_id = SYSTEM_USER_ID()
);
-- 3. Audit sensitive operations
CREATE SERVER AUDIT DatabaseAudit
TO FILE (FILEPATH = 'C:\Audit\');
CREATE DATABASE AUDIT SPECIFICATION DatabaseAuditSpec
FOR SERVER AUDIT DatabaseAudit
ADD (SELECT, INSERT, UPDATE, DELETE ON DATABASE::MyDatabase BY public),
ADD (EXECUTE ON DATABASE::MyDatabase BY public);
-- 4. Implement data masking
CREATE TABLE customer_data (
customer_id INT PRIMARY KEY,
full_name VARCHAR(100) MASKED WITH (FUNCTION = 'partial(2, "XXX", 2)'),
email VARCHAR(100) MASKED WITH (FUNCTION = 'email()'),
phone VARCHAR(15) MASKED WITH (FUNCTION = 'partial(0, "XXX-XXX-", 4)'),
credit_card VARCHAR(20) MASKED WITH (FUNCTION = 'partial(0, "XXXX-XXXX-XXXX-", 4)')
);
-- 5. Create secure connection requirements
CREATE LOGIN secure_user WITH PASSWORD = 'ComplexPassword123!';
CREATE USER secure_user FOR LOGIN secure_user;
-- Require SSL connections
ALTER LOGIN secure_user
WITH PASSWORD = 'ComplexPassword123!'
CHECK_POLICY = ON,
CHECK_EXPIRATION = ON;
-- Grant minimal required permissions
GRANT SELECT, INSERT ON customer_data TO secure_user;
DENY DELETE, UPDATE ON customer_data TO secure_user;
91. Database Data Quality Checks
-- 1. Data completeness checks
CREATE PROCEDURE CheckDataCompleteness
AS
BEGIN
DECLARE @completeness_score DECIMAL(5,2);
SELECT @completeness_score =
(COUNT(*) * 100.0) /
(SELECT COUNT(*) FROM customer_data)
FROM customer_data
WHERE full_name IS NOT NULL
AND email IS NOT NULL
AND phone IS NOT NULL;
INSERT INTO data_quality_log (
check_type,
score,
check_date,
table_name
) VALUES (
'completeness',
@completeness_score,
GETDATE(),
'customer_data'
);
-- Alert if score is below threshold
IF @completeness_score < 95.0
EXEC SendDataQualityAlert 'Completeness check failed', @completeness_score;
END;
-- 2. Data accuracy validation
CREATE FUNCTION ValidateEmailFormat(@email VARCHAR(100))
RETURNS BIT
AS
BEGIN
DECLARE @is_valid BIT = 0;
IF @email LIKE '%_@_%._%'
AND @email NOT LIKE '%@%@%'
AND @email NOT LIKE '%..%'
AND @email NOT LIKE '%.@%'
AND @email NOT LIKE '%@.%'
SET @is_valid = 1;
RETURN @is_valid;
END;
-- 3. Data consistency checks
CREATE PROCEDURE CheckReferentialIntegrity
AS
BEGIN
-- Check for orphaned records
SELECT 'Orphaned orders found' as issue_type, COUNT(*) as count
FROM orders o
LEFT JOIN customers c ON o.customer_id = c.customer_id
WHERE c.customer_id IS NULL;
-- Check for duplicate records
SELECT 'Duplicate emails found' as issue_type, COUNT(*) as count
FROM (
SELECT email, COUNT(*) as cnt
FROM customer_data
GROUP BY email
HAVING COUNT(*) > 1
) duplicates;
END;
-- 4. Data range validation
CREATE TRIGGER ValidateSalaryRange
ON employee_data
AFTER INSERT, UPDATE
AS
BEGIN
IF EXISTS (
SELECT 1 FROM inserted
WHERE salary < 0 OR salary > 1000000
)
BEGIN
RAISERROR ('Salary must be between 0 and 1,000,000', 16, 1);
ROLLBACK TRANSACTION;
END;
END;
92. Database Data Lineage Tracking
-- 1. Create lineage tracking tables
CREATE TABLE data_lineage (
lineage_id INT IDENTITY(1,1) PRIMARY KEY,
source_table VARCHAR(100),
source_column VARCHAR(100),
target_table VARCHAR(100),
target_column VARCHAR(100),
transformation_type VARCHAR(50),
transformation_sql TEXT,
created_date DATETIME2 DEFAULT GETDATE(),
created_by VARCHAR(50) DEFAULT SYSTEM_USER
);
CREATE TABLE data_flow_log (
flow_id INT IDENTITY(1,1) PRIMARY KEY,
source_system VARCHAR(100),
target_system VARCHAR(100),
data_volume INT,
processing_time_ms INT,
status VARCHAR(20),
error_message TEXT,
flow_date DATETIME2 DEFAULT GETDATE()
);
-- 2. Track ETL transformations
CREATE PROCEDURE TrackDataTransformation
@source_table VARCHAR(100),
@target_table VARCHAR(100),
@transformation_sql TEXT
AS
BEGIN
INSERT INTO data_lineage (
source_table,
target_table,
transformation_type,
transformation_sql
) VALUES (
@source_table,
@target_table,
'ETL_TRANSFORMATION',
@transformation_sql
);
-- Log the actual transformation
DECLARE @start_time DATETIME2 = GETDATE();
DECLARE @row_count INT;
BEGIN TRY
EXEC sp_executesql @transformation_sql;
SELECT @row_count = @@ROWCOUNT;
INSERT INTO data_flow_log (
source_system,
target_system,
data_volume,
processing_time_ms,
status
) VALUES (
@source_table,
@target_table,
@row_count,
DATEDIFF(MILLISECOND, @start_time, GETDATE()),
'SUCCESS'
);
END TRY
BEGIN CATCH
INSERT INTO data_flow_log (
source_system,
target_system,
status,
error_message
) VALUES (
@source_table,
@target_table,
'FAILED',
ERROR_MESSAGE()
);
END CATCH
END;
-- 3. Query lineage information
CREATE VIEW v_data_lineage_report AS
SELECT
source_table,
target_table,
transformation_type,
COUNT(*) as transformation_count,
MAX(created_date) as last_transformation
FROM data_lineage
GROUP BY source_table, target_table, transformation_type;
-- 4. Track data dependencies
CREATE PROCEDURE AnalyzeDataDependencies
@table_name VARCHAR(100)
AS
BEGIN
-- Find all tables that depend on this table
SELECT
target_table as dependent_table,
transformation_type,
created_date
FROM data_lineage
WHERE source_table = @table_name
ORDER BY created_date DESC;
-- Find all tables this table depends on
SELECT
source_table as dependency_table,
transformation_type,
created_date
FROM data_lineage
WHERE target_table = @table_name
ORDER BY created_date DESC;
END;
93. Database Data Governance Policies
-- 1. Create governance framework tables
CREATE TABLE governance_policies (
policy_id INT IDENTITY(1,1) PRIMARY KEY,
policy_name VARCHAR(100),
policy_type VARCHAR(50), -- ACCESS, RETENTION, QUALITY, etc.
policy_description TEXT,
is_active BIT DEFAULT 1,
created_date DATETIME2 DEFAULT GETDATE(),
created_by VARCHAR(50)
);
CREATE TABLE policy_applications (
application_id INT IDENTITY(1,1) PRIMARY KEY,
policy_id INT,
table_name VARCHAR(100),
column_name VARCHAR(100),
application_date DATETIME2 DEFAULT GETDATE(),
FOREIGN KEY (policy_id) REFERENCES governance_policies(policy_id)
);
-- 2. Implement access governance
CREATE PROCEDURE ApplyAccessGovernance
@table_name VARCHAR(100),
@user_role VARCHAR(50)
AS
BEGIN
DECLARE @access_level VARCHAR(20);
-- Determine access level based on role
SELECT @access_level =
CASE @user_role
WHEN 'ADMIN' THEN 'FULL'
WHEN 'MANAGER' THEN 'READ_WRITE'
WHEN 'ANALYST' THEN 'READ_ONLY'
ELSE 'RESTRICTED'
END;
-- Apply appropriate permissions
IF @access_level = 'FULL'
GRANT ALL ON @table_name TO @user_role;
ELSE IF @access_level = 'READ_WRITE'
GRANT SELECT, INSERT, UPDATE ON @table_name TO @user_role;
ELSE IF @access_level = 'READ_ONLY'
GRANT SELECT ON @table_name TO @user_role;
ELSE
DENY ALL ON @table_name TO @user_role;
-- Log the policy application
INSERT INTO policy_applications (policy_id, table_name)
SELECT policy_id, @table_name
FROM governance_policies
WHERE policy_type = 'ACCESS' AND is_active = 1;
END;
-- 3. Data quality governance
CREATE PROCEDURE EnforceDataQualityGovernance
@table_name VARCHAR(100)
AS
BEGIN
-- Check for required fields
DECLARE @sql NVARCHAR(MAX) =
'SELECT COUNT(*) FROM ' + @table_name +
' WHERE required_field IS NULL OR required_field = ''''';
DECLARE @null_count INT;
EXEC sp_executesql @sql, N'@count INT OUTPUT', @count = @null_count OUTPUT;
IF @null_count > 0
BEGIN
-- Log violation
INSERT INTO governance_violations (
table_name,
violation_type,
violation_count,
violation_date
) VALUES (
@table_name,
'NULL_REQUIRED_FIELD',
@null_count,
GETDATE()
);
-- Send alert
EXEC SendGovernanceAlert @table_name, 'Data quality violation detected';
END;
END;
-- 4. Compliance monitoring
CREATE VIEW v_compliance_status AS
SELECT
p.policy_name,
p.policy_type,
COUNT(pa.application_id) as applications,
CASE
WHEN p.is_active = 1 THEN 'ACTIVE'
ELSE 'INACTIVE'
END as status
FROM governance_policies p
LEFT JOIN policy_applications pa ON p.policy_id = pa.policy_id
GROUP BY p.policy_id, p.policy_name, p.policy_type, p.is_active;
94. Database Data Cataloging
-- 1. Create data catalog tables
CREATE TABLE data_catalog (
catalog_id INT IDENTITY(1,1) PRIMARY KEY,
table_name VARCHAR(100),
column_name VARCHAR(100),
data_type VARCHAR(50),
description TEXT,
business_owner VARCHAR(100),
technical_owner VARCHAR(100),
data_classification VARCHAR(50), -- PUBLIC, INTERNAL, CONFIDENTIAL, RESTRICTED
last_updated DATETIME2 DEFAULT GETDATE(),
created_date DATETIME2 DEFAULT GETDATE()
);
CREATE TABLE data_dictionary (
dictionary_id INT IDENTITY(1,1) PRIMARY KEY,
term_name VARCHAR(100),
definition TEXT,
business_context TEXT,
technical_context TEXT,
related_terms VARCHAR(500),
created_date DATETIME2 DEFAULT GETDATE()
);
-- 2. Auto-generate catalog entries
CREATE PROCEDURE GenerateDataCatalog
@database_name VARCHAR(100)
AS
BEGIN
DECLARE @sql NVARCHAR(MAX) = '
INSERT INTO data_catalog (table_name, column_name, data_type)
SELECT
t.name as table_name,
c.name as column_name,
ty.name +
CASE
WHEN c.max_length = -1 THEN ''(MAX)''
WHEN ty.name IN (''varchar'', ''char'', ''nvarchar'', ''nchar'')
THEN ''('' + CAST(c.max_length/2 AS VARCHAR) + '')''
WHEN ty.name IN (''decimal'', ''numeric'')
THEN ''('' + CAST(c.precision AS VARCHAR) + '','' + CAST(c.scale AS VARCHAR) + '')''
ELSE ''''
END as data_type
FROM ' + @database_name + '.sys.tables t
INNER JOIN ' + @database_name + '.sys.columns c ON t.object_id = c.object_id
INNER JOIN ' + @database_name + '.sys.types ty ON c.user_type_id = ty.user_type_id
WHERE t.is_ms_shipped = 0
ORDER BY t.name, c.column_id';
EXEC sp_executesql @sql;
END;
-- 3. Search and discovery
CREATE PROCEDURE SearchDataCatalog
@search_term VARCHAR(100)
AS
BEGIN
SELECT
table_name,
column_name,
data_type,
description,
business_owner,
data_classification
FROM data_catalog
WHERE table_name LIKE '%' + @search_term + '%'
OR column_name LIKE '%' + @search_term + '%'
OR description LIKE '%' + @search_term + '%'
OR business_owner LIKE '%' + @search_term + '%'
ORDER BY table_name, column_name;
END;
-- 4. Data lineage in catalog
CREATE VIEW v_catalog_with_lineage AS
SELECT
c.*,
l.source_table,
l.transformation_type,
l.created_date as last_transformation
FROM data_catalog c
LEFT JOIN data_lineage l ON c.table_name = l.target_table
AND c.column_name = l.target_column;
95. Database Data Profiling
-- 1. Create profiling tables
CREATE TABLE data_profile_results (
profile_id INT IDENTITY(1,1) PRIMARY KEY,
table_name VARCHAR(100),
column_name VARCHAR(100),
total_rows INT,
null_count INT,
distinct_count INT,
min_value VARCHAR(500),
max_value VARCHAR(500),
avg_length DECIMAL(10,2),
data_type VARCHAR(50),
profile_date DATETIME2 DEFAULT GETDATE()
);
-- 2. Comprehensive data profiling procedure
CREATE PROCEDURE ProfileTableData
@table_name VARCHAR(100)
AS
BEGIN
DECLARE @sql NVARCHAR(MAX);
DECLARE @column_name VARCHAR(100);
DECLARE @data_type VARCHAR(50);
-- Get all columns for the table
DECLARE column_cursor CURSOR FOR
SELECT name, system_type_name
FROM sys.columns c
INNER JOIN sys.types t ON c.user_type_id = t.user_type_id
WHERE object_id = OBJECT_ID(@table_name);
OPEN column_cursor;
FETCH NEXT FROM column_cursor INTO @column_name, @data_type;
WHILE @@FETCH_STATUS = 0
BEGIN
-- Build dynamic SQL for profiling
SET @sql = '
INSERT INTO data_profile_results (
table_name, column_name, total_rows, null_count,
distinct_count, min_value, max_value, avg_length, data_type
)
SELECT
''' + @table_name + ''' as table_name,
''' + @column_name + ''' as column_name,
COUNT(*) as total_rows,
SUM(CASE WHEN [' + @column_name + '] IS NULL THEN 1 ELSE 0 END) as null_count,
COUNT(DISTINCT [' + @column_name + ']) as distinct_count,
MIN(CAST([' + @column_name + '] AS VARCHAR(500))) as min_value,
MAX(CAST([' + @column_name + '] AS VARCHAR(500))) as max_value,
AVG(LEN(CAST([' + @column_name + '] AS VARCHAR(500)))) as avg_length,
''' + @data_type + ''' as data_type
FROM ' + @table_name;
EXEC sp_executesql @sql;
FETCH NEXT FROM column_cursor INTO @column_name, @data_type;
END;
CLOSE column_cursor;
DEALLOCATE column_cursor;
END;
-- 3. Pattern analysis
CREATE PROCEDURE AnalyzeDataPatterns
@table_name VARCHAR(100),
@column_name VARCHAR(100)
AS
BEGIN
DECLARE @sql NVARCHAR(MAX) = '
WITH pattern_analysis AS (
SELECT
[' + @column_name + '],
LEN([' + @column_name + ']) as length,
PATINDEX(''%[^0-9]%'', [' + @column_name + ']) as non_numeric_pos,
PATINDEX(''%[^A-Za-z]%'', [' + @column_name + ']) as non_alpha_pos
FROM ' + @table_name + '
WHERE [' + @column_name + '] IS NOT NULL
)
SELECT
COUNT(*) as total_values,
COUNT(CASE WHEN non_numeric_pos = 0 THEN 1 END) as numeric_only,
COUNT(CASE WHEN non_alpha_pos = 0 THEN 1 END) as alpha_only,
COUNT(CASE WHEN non_numeric_pos > 0 AND non_alpha_pos > 0 THEN 1 END) as mixed,
AVG(length) as avg_length,
MIN(length) as min_length,
MAX(length) as max_length
FROM pattern_analysis';
EXEC sp_executesql @sql;
END;
-- 4. Data quality score calculation
CREATE FUNCTION CalculateDataQualityScore(@table_name VARCHAR(100))
RETURNS DECIMAL(5,2)
AS
BEGIN
DECLARE @quality_score DECIMAL(5,2);
SELECT @quality_score =
AVG(
CASE
WHEN null_count = 0 THEN 100.0
ELSE ((total_rows - null_count) * 100.0) / total_rows
END
)
FROM data_profile_results
WHERE table_name = @table_name;
RETURN ISNULL(@quality_score, 0);
END;
96. Database Data Stewardship
-- 1. Create stewardship framework
CREATE TABLE data_stewards (
steward_id INT IDENTITY(1,1) PRIMARY KEY,
steward_name VARCHAR(100),
email VARCHAR(100),
department VARCHAR(100),
assigned_domains VARCHAR(500),
is_active BIT DEFAULT 1,
created_date DATETIME2 DEFAULT GETDATE()
);
CREATE TABLE stewardship_assignments (
assignment_id INT IDENTITY(1,1) PRIMARY KEY,
steward_id INT,
table_name VARCHAR(100),
column_name VARCHAR(100),
assignment_type VARCHAR(50), -- OWNER, REVIEWER, APPROVER
assignment_date DATETIME2 DEFAULT GETDATE(),
FOREIGN KEY (steward_id) REFERENCES data_stewards(steward_id)
);
-- 2. Stewardship approval workflow
CREATE TABLE stewardship_approvals (
approval_id INT IDENTITY(1,1) PRIMARY KEY,
change_request_id INT,
steward_id INT,
approval_status VARCHAR(20), -- PENDING, APPROVED, REJECTED
approval_notes TEXT,
approval_date DATETIME2,
FOREIGN KEY (steward_id) REFERENCES data_stewards(steward_id)
);
CREATE PROCEDURE RequestDataChange
@table_name VARCHAR(100),
@change_description TEXT,
@requested_by VARCHAR(100)
AS
BEGIN
DECLARE @change_request_id INT;
-- Create change request
INSERT INTO data_change_requests (
table_name,
change_description,
requested_by,
request_date
) VALUES (
@table_name,
@change_description,
@requested_by,
GETDATE()
);
SET @change_request_id = SCOPE_IDENTITY();
-- Notify relevant stewards
INSERT INTO stewardship_approvals (change_request_id, steward_id, approval_status)
SELECT @change_request_id, steward_id, 'PENDING'
FROM stewardship_assignments sa
INNER JOIN data_stewards ds ON sa.steward_id = ds.steward_id
WHERE sa.table_name = @table_name
AND ds.is_active = 1;
-- Send notification emails
EXEC SendStewardNotification @change_request_id;
END;
-- 3. Data quality monitoring by stewards
CREATE PROCEDURE MonitorDataQualityBySteward
@steward_id INT
AS
BEGIN
SELECT
sa.table_name,
sa.column_name,
dpr.total_rows,
dpr.null_count,
dpr.distinct_count,
((dpr.total_rows - dpr.null_count) * 100.0) / dpr.total_rows as completeness_score
FROM stewardship_assignments sa
INNER JOIN data_profile_results dpr ON sa.table_name = dpr.table_name
AND sa.column_name = dpr.column_name
WHERE sa.steward_id = @steward_id
AND dpr.profile_date = (
SELECT MAX(profile_date)
FROM data_profile_results
WHERE table_name = sa.table_name
AND column_name = sa.column_name
)
ORDER BY completeness_score ASC;
END;
-- 4. Stewardship reporting
CREATE VIEW v_stewardship_summary AS
SELECT
ds.steward_name,
ds.department,
COUNT(sa.assignment_id) as assigned_assets,
COUNT(DISTINCT sa.table_name) as assigned_tables,
AVG(dq.quality_score) as avg_quality_score
FROM data_stewards ds
LEFT JOIN stewardship_assignments sa ON ds.steward_id = sa.steward_id
LEFT JOIN (
SELECT
table_name,
AVG(
CASE
WHEN null_count = 0 THEN 100.0
ELSE ((total_rows - null_count) * 100.0) / total_rows
END
) as quality_score
FROM data_profile_results
GROUP BY table_name
) dq ON sa.table_name = dq.table_name
WHERE ds.is_active = 1
GROUP BY ds.steward_id, ds.steward_name, ds.department;
97. Database Data Classification
-- 1. Create classification framework
CREATE TABLE data_classifications (
classification_id INT IDENTITY(1,1) PRIMARY KEY,
classification_level VARCHAR(50), -- PUBLIC, INTERNAL, CONFIDENTIAL, RESTRICTED
description TEXT,
retention_period_months INT,
encryption_required BIT DEFAULT 0,
audit_required BIT DEFAULT 0,
created_date DATETIME2 DEFAULT GETDATE()
);
CREATE TABLE classified_assets (
asset_id INT IDENTITY(1,1) PRIMARY KEY,
table_name VARCHAR(100),
column_name VARCHAR(100),
classification_id INT,
classification_reason TEXT,
classified_by VARCHAR(100),
classification_date DATETIME2 DEFAULT GETDATE(),
FOREIGN KEY (classification_id) REFERENCES data_classifications(classification_id)
);
-- 2. Auto-classification based on patterns
CREATE PROCEDURE AutoClassifyData
@table_name VARCHAR(100)
AS
BEGIN
DECLARE @sql NVARCHAR(MAX);
DECLARE @column_name VARCHAR(100);
-- Keywords for classification
DECLARE @confidential_keywords TABLE (keyword VARCHAR(50));
INSERT INTO @confidential_keywords VALUES
('password'), ('ssn'), ('credit_card'), ('salary'),
('address'), ('phone'), ('email'), ('dob');
DECLARE column_cursor CURSOR FOR
SELECT name FROM sys.columns
WHERE object_id = OBJECT_ID(@table_name);
OPEN column_cursor;
FETCH NEXT FROM column_cursor INTO @column_name;
WHILE @@FETCH_STATUS = 0
BEGIN
DECLARE @classification_id INT;
-- Determine classification based on column name
SELECT @classification_id =
CASE
WHEN EXISTS (
SELECT 1 FROM @confidential_keywords
WHERE LOWER(@column_name) LIKE '%' + keyword + '%'
) THEN 3 -- CONFIDENTIAL
WHEN @column_name LIKE '%id%' THEN 2 -- INTERNAL
ELSE 1 -- PUBLIC
END;
-- Apply classification
INSERT INTO classified_assets (
table_name,
column_name,
classification_id,
classified_by
) VALUES (
@table_name,
@column_name,
@classification_id,
SYSTEM_USER
);
FETCH NEXT FROM column_cursor INTO @column_name;
END;
CLOSE column_cursor;
DEALLOCATE column_cursor;
END;
-- 3. Classification-based access control
CREATE PROCEDURE ApplyClassificationBasedAccess
@user_id VARCHAR(100),
@user_clearance_level VARCHAR(50)
AS
BEGIN
DECLARE @max_classification_id INT;
-- Determine user's maximum allowed classification
SELECT @max_classification_id =
CASE @user_clearance_level
WHEN 'PUBLIC' THEN 1
WHEN 'INTERNAL' THEN 2
WHEN 'CONFIDENTIAL' THEN 3
WHEN 'RESTRICTED' THEN 4
ELSE 1
END;
-- Grant access to tables based on classification
DECLARE @sql NVARCHAR(MAX) = '
GRANT SELECT ON ' +
(SELECT STRING_AGG(table_name, ', ')
FROM classified_assets ca
INNER JOIN data_classifications dc ON ca.classification_id = dc.classification_id
WHERE ca.classification_id <= @max_classification_id
GROUP BY ca.table_name) +
' TO [' + @user_id + ']';
EXEC sp_executesql @sql;
END;
-- 4. Classification compliance reporting
CREATE VIEW v_classification_compliance AS
SELECT
dc.classification_level,
COUNT(ca.asset_id) as asset_count,
COUNT(CASE WHEN dc.encryption_required = 1 THEN 1 END) as encrypted_assets,
COUNT(CASE WHEN dc.audit_required = 1 THEN 1 END) as audited_assets,
AVG(DATEDIFF(MONTH, ca.classification_date, GETDATE())) as avg_age_months
FROM data_classifications dc
LEFT JOIN classified_assets ca ON dc.classification_id = ca.classification_id
GROUP BY dc.classification_id, dc.classification_level;
98. Database Data Retention Policies
-- 1. Create retention policy framework
CREATE TABLE retention_policies (
policy_id INT IDENTITY(1,1) PRIMARY KEY,
policy_name VARCHAR(100),
table_name VARCHAR(100),
retention_period_months INT,
retention_criteria VARCHAR(500), -- Column to base retention on
archive_required BIT DEFAULT 0,
delete_required BIT DEFAULT 1,
is_active BIT DEFAULT 1,
created_date DATETIME2 DEFAULT GETDATE()
);
CREATE TABLE retention_execution_log (
execution_id INT IDENTITY(1,1) PRIMARY KEY,
policy_id INT,
records_processed INT,
records_archived INT,
records_deleted INT,
execution_date DATETIME2 DEFAULT GETDATE(),
execution_status VARCHAR(20), -- SUCCESS, FAILED, PARTIAL
error_message TEXT,
FOREIGN KEY (policy_id) REFERENCES retention_policies(policy_id)
);
-- 2. Automated retention enforcement
CREATE PROCEDURE EnforceRetentionPolicies
AS
BEGIN
DECLARE @policy_id INT;
DECLARE @table_name VARCHAR(100);
DECLARE @retention_period_months INT;
DECLARE @retention_criteria VARCHAR(500);
DECLARE @archive_required BIT;
DECLARE @delete_required BIT;
DECLARE @sql NVARCHAR(MAX);
DECLARE @cutoff_date DATE;
DECLARE policy_cursor CURSOR FOR
SELECT
policy_id,
table_name,
retention_period_months,
retention_criteria,
archive_required,
delete_required
FROM retention_policies
WHERE is_active = 1;
OPEN policy_cursor;
FETCH NEXT FROM policy_cursor INTO
@policy_id, @table_name, @retention_period_months,
@retention_criteria, @archive_required, @delete_required;
WHILE @@FETCH_STATUS = 0
BEGIN
SET @cutoff_date = DATEADD(MONTH, -@retention_period_months, GETDATE());
BEGIN TRY
-- Archive if required
IF @archive_required = 1
BEGIN
SET @sql = '
INSERT INTO ' + @table_name + '_archive
SELECT * FROM ' + @table_name + '
WHERE ' + @retention_criteria + ' < ''' + CAST(@cutoff_date AS VARCHAR) + '''';
EXEC sp_executesql @sql;
SET @sql = 'SELECT @count = @@ROWCOUNT';
DECLARE @archived_count INT;
EXEC sp_executesql @sql, N'@count INT OUTPUT', @count = @archived_count OUTPUT;
END;
-- Delete if required
IF @delete_required = 1
BEGIN
SET @sql = '
DELETE FROM ' + @table_name + '
WHERE ' + @retention_criteria + ' < ''' + CAST(@cutoff_date AS VARCHAR) + '''';
EXEC sp_executesql @sql;
SET @sql = 'SELECT @count = @@ROWCOUNT';
DECLARE @deleted_count INT;
EXEC sp_executesql @sql, N'@count INT OUTPUT', @count = @deleted_count OUTPUT;
END;
-- Log successful execution
INSERT INTO retention_execution_log (
policy_id,
records_archived,
records_deleted,
execution_status
) VALUES (
@policy_id,
ISNULL(@archived_count, 0),
ISNULL(@deleted_count, 0),
'SUCCESS'
);
END TRY
BEGIN CATCH
-- Log failed execution
INSERT INTO retention_execution_log (
policy_id,
execution_status,
error_message
) VALUES (
@policy_id,
'FAILED',
ERROR_MESSAGE()
);
END CATCH
FETCH NEXT FROM policy_cursor INTO
@policy_id, @table_name, @retention_period_months,
@retention_criteria, @archive_required, @delete_required;
END;
CLOSE policy_cursor;
DEALLOCATE policy_cursor;
END;
-- 3. Retention compliance monitoring
CREATE VIEW v_retention_compliance AS
SELECT
rp.policy_name,
rp.table_name,
rp.retention_period_months,
COUNT(rel.execution_id) as execution_count,
MAX(rel.execution_date) as last_execution,
CASE
WHEN MAX(rel.execution_date) < DATEADD(DAY, -7, GETDATE()) THEN 'OVERDUE'
WHEN MAX(rel.execution_date) < DATEADD(DAY, -1, GETDATE()) THEN 'DUE'
ELSE 'CURRENT'
END as compliance_status
FROM retention_policies rp
LEFT JOIN retention_execution_log rel ON rp.policy_id = rel.policy_id
WHERE rp.is_active = 1
GROUP BY rp.policy_id, rp.policy_name, rp.table_name, rp.retention_period_months;
-- 4. Legal hold management
CREATE TABLE legal_holds (
hold_id INT IDENTITY(1,1) PRIMARY KEY,
case_number VARCHAR(100),
table_name VARCHAR(100),
hold_reason TEXT,
hold_start_date DATETIME2 DEFAULT GETDATE(),
hold_end_date DATETIME2,
is_active BIT DEFAULT 1,
created_by VARCHAR(100)
);
CREATE PROCEDURE ApplyLegalHold
@case_number VARCHAR(100),
@table_name VARCHAR(100),
@hold_reason TEXT
AS
BEGIN
-- Create legal hold
INSERT INTO legal_holds (
case_number,
table_name,
hold_reason,
created_by
) VALUES (
@case_number,
@table_name,
@hold_reason,
SYSTEM_USER
);
-- Disable retention policies for this table
UPDATE retention_policies
SET is_active = 0
WHERE table_name = @table_name;
-- Log the action
INSERT INTO retention_execution_log (
policy_id,
execution_status,
error_message
) VALUES (
NULL,
'LEGAL_HOLD',
'Legal hold applied for case: ' + @case_number
);
END;
99. Database Data Privacy Compliance
-- 1. Create privacy compliance framework
CREATE TABLE privacy_regulations (
regulation_id INT IDENTITY(1,1) PRIMARY KEY,
regulation_name VARCHAR(100), -- GDPR, CCPA, HIPAA, etc.
regulation_description TEXT,
effective_date DATE,
is_active BIT DEFAULT 1
);
CREATE TABLE personal_data_fields (
field_id INT IDENTITY(1,1) PRIMARY KEY,
table_name VARCHAR(100),
column_name VARCHAR(100),
data_category VARCHAR(50), -- PII, PHI, FINANCIAL, etc.
sensitivity_level VARCHAR(20), -- LOW, MEDIUM, HIGH, CRITICAL
regulation_applicable VARCHAR(500), -- Comma-separated regulation IDs
created_date DATETIME2 DEFAULT GETDATE()
);
-- 2. GDPR compliance implementation
CREATE PROCEDURE ImplementGDPRCompliance
AS
BEGIN
-- Right to be forgotten (Data deletion)
CREATE PROCEDURE DeletePersonalData
@user_identifier VARCHAR(100),
@identifier_type VARCHAR(50) -- email, user_id, etc.
AS
BEGIN
DECLARE @sql NVARCHAR(MAX);
-- Anonymize or delete personal data
UPDATE users
SET
email = 'deleted_' + CAST(user_id AS VARCHAR) + '@deleted.com',
full_name = 'DELETED_USER',
phone = NULL,
address = NULL
WHERE @identifier_type = 'user_id' AND user_id = @user_identifier
OR @identifier_type = 'email' AND email = @user_identifier;
-- Log the deletion
INSERT INTO data_deletion_log (
user_identifier,
identifier_type,
deletion_date,
deleted_by
) VALUES (
@user_identifier,
@identifier_type,
GETDATE(),
SYSTEM_USER
);
END;
-- Data portability (Export personal data)
CREATE PROCEDURE ExportPersonalData
@user_identifier VARCHAR(100),
@identifier_type VARCHAR(50)
AS
BEGIN
SELECT
user_id,
email,
full_name,
phone,
address,
created_date,
last_login_date
FROM users
WHERE @identifier_type = 'user_id' AND user_id = @user_identifier
OR @identifier_type = 'email' AND email = @user_identifier;
END;
END;
-- 3. Data anonymization for privacy
CREATE PROCEDURE AnonymizePersonalData
@table_name VARCHAR(100),
@anonymization_level VARCHAR(20) -- FULL, PARTIAL, HASH
AS
BEGIN
DECLARE @sql NVARCHAR(MAX);
IF @anonymization_level = 'FULL'
BEGIN
SET @sql = '
UPDATE ' + @table_name + '
SET
email = ''anonymous@example.com'',
full_name = ''ANONYMOUS_USER'',
phone = ''000-000-0000'',
address = ''ANONYMOUS_ADDRESS''';
END
ELSE IF @anonymization_level = 'HASH'
BEGIN
SET @sql = '
UPDATE ' + @table_name + '
SET
email = HASHBYTES(''SHA2_256'', email),
full_name = HASHBYTES(''SHA2_256'', full_name),
phone = HASHBYTES(''SHA2_256'', phone)';
END;
EXEC sp_executesql @sql;
-- Log anonymization
INSERT INTO privacy_audit_log (
table_name,
operation_type,
operation_date,
performed_by
) VALUES (
@table_name,
'ANONYMIZATION',
GETDATE(),
SYSTEM_USER
);
END;
-- 4. Privacy impact assessment
CREATE VIEW v_privacy_impact_assessment AS
SELECT
pdf.table_name,
pdf.column_name,
pdf.data_category,
pdf.sensitivity_level,
COUNT(*) as record_count,
pr.regulation_name,
CASE
WHEN pdf.sensitivity_level IN ('HIGH', 'CRITICAL') THEN 'HIGH_RISK'
WHEN pdf.sensitivity_level = 'MEDIUM' THEN 'MEDIUM_RISK'
ELSE 'LOW_RISK'
END as privacy_risk_level
FROM personal_data_fields pdf
INNER JOIN privacy_regulations pr ON
CHARINDEX(CAST(pr.regulation_id AS VARCHAR), pdf.regulation_applicable) > 0
WHERE pr.is_active = 1
GROUP BY pdf.table_name, pdf.column_name, pdf.data_category,
pdf.sensitivity_level, pr.regulation_name;
Database Data Ethics Considerations - Technical Implementation
1. Data Privacy and Anonymization
Pseudonymization with SQL:
-- Create a mapping table for pseudonymization
CREATE TABLE user_pseudonym_mapping (
original_id INT PRIMARY KEY,
pseudonym_id VARCHAR(50) UNIQUE,
created_date DATETIME DEFAULT GETDATE(),
expiry_date DATETIME
);
-- Function to generate pseudonyms
CREATE FUNCTION GeneratePseudonym(@originalId INT)
RETURNS VARCHAR(50)
AS
BEGIN
DECLARE @pseudonym VARCHAR(50)
SET @pseudonym = CONCAT('USER_', CAST(@originalId AS VARCHAR), '_',
CAST(ABS(CHECKSUM(NEWID())) % 10000 AS VARCHAR))
RETURN @pseudonym
END;
-- View for anonymized data access
CREATE VIEW anonymized_user_data AS
SELECT
upm.pseudonym_id,
ud.email_domain, -- Only domain, not full email
ud.age_group, -- Age ranges instead of exact age
ud.region -- Broader geographic area
FROM user_data ud
JOIN user_pseudonym_mapping upm ON ud.user_id = upm.original_id
WHERE upm.expiry_date > GETDATE();
2. Data Retention and Lifecycle Management
Automated Data Retention Policy:
-- Create retention policy table
CREATE TABLE data_retention_policy (
table_name VARCHAR(100),
retention_period_months INT,
archive_before_delete BIT DEFAULT 1,
last_cleanup_date DATETIME
);
-- Stored procedure for automated cleanup
CREATE PROCEDURE ExecuteDataRetentionPolicy
@tableName VARCHAR(100)
AS
BEGIN
SET NOCOUNT ON;
DECLARE @retentionMonths INT, @cutoffDate DATETIME;
SELECT @retentionMonths = retention_period_months
FROM data_retention_policy
WHERE table_name = @tableName;
SET @cutoffDate = DATEADD(MONTH, -@retentionMonths, GETDATE());
-- Archive data before deletion
IF EXISTS (SELECT 1 FROM data_retention_policy
WHERE table_name = @tableName AND archive_before_delete = 1)
BEGIN
EXEC ArchiveTableData @tableName, @cutoffDate;
END
-- Delete expired data
DECLARE @sql NVARCHAR(MAX) =
'DELETE FROM ' + @tableName +
' WHERE created_date < @cutoffDate';
EXEC sp_executesql @sql, N'@cutoffDate DATETIME', @cutoffDate;
-- Update last cleanup date
UPDATE data_retention_policy
SET last_cleanup_date = GETDATE()
WHERE table_name = @tableName;
END;
3. Data Access Control and Audit Trail
Comprehensive Audit System:
-- Audit trail table
CREATE TABLE data_access_audit (
audit_id BIGINT IDENTITY(1,1) PRIMARY KEY,
user_id VARCHAR(50),
action_type VARCHAR(20), -- SELECT, INSERT, UPDATE, DELETE
table_name VARCHAR(100),
record_id VARCHAR(50),
old_values NVARCHAR(MAX),
new_values NVARCHAR(MAX),
ip_address VARCHAR(45),
session_id VARCHAR(100),
access_timestamp DATETIME DEFAULT GETDATE(),
purpose VARCHAR(200),
data_classification VARCHAR(20) -- PUBLIC, INTERNAL, CONFIDENTIAL, RESTRICTED
);
-- Trigger for automatic audit logging
CREATE TRIGGER tr_audit_user_data_changes
ON user_data
AFTER INSERT, UPDATE, DELETE
AS
BEGIN
DECLARE @action VARCHAR(20), @oldValues NVARCHAR(MAX), @newValues NVARCHAR(MAX);
IF EXISTS(SELECT 1 FROM inserted) AND EXISTS(SELECT 1 FROM deleted)
SET @action = 'UPDATE';
ELSE IF EXISTS(SELECT 1 FROM inserted)
SET @action = 'INSERT';
ELSE
SET @action = 'DELETE';
-- Capture old and new values for sensitive fields
SELECT @oldValues = (
SELECT user_id, email, phone, ssn
FROM deleted
FOR JSON PATH
);
SELECT @newValues = (
SELECT user_id, email, phone, ssn
FROM inserted
FOR JSON PATH
);
INSERT INTO data_access_audit (
user_id, action_type, table_name, record_id,
old_values, new_values, ip_address, session_id, purpose
)
VALUES (
SYSTEM_USER, @action, 'user_data',
COALESCE((SELECT user_id FROM inserted), (SELECT user_id FROM deleted)),
@oldValues, @newValues,
CAST(CONNECTIONPROPERTY('client_net_address') AS VARCHAR(45)),
@@SPID, 'Data maintenance'
);
END;
4. Data Classification and Handling
Data Classification System:
-- Data classification metadata
CREATE TABLE data_classification (
table_name VARCHAR(100),
column_name VARCHAR(100),
classification_level VARCHAR(20), -- PUBLIC, INTERNAL, CONFIDENTIAL, RESTRICTED
pii_flag BIT DEFAULT 0, -- Personally Identifiable Information
phi_flag BIT DEFAULT 0, -- Protected Health Information
pci_flag BIT DEFAULT 0, -- Payment Card Industry data
encryption_required BIT DEFAULT 0,
masking_required BIT DEFAULT 0,
access_approval_required BIT DEFAULT 0
);
-- Function to check data access permissions
CREATE FUNCTION CheckDataAccessPermission(
@userId VARCHAR(50),
@tableName VARCHAR(100),
@classificationLevel VARCHAR(20)
)
RETURNS BIT
AS
BEGIN
DECLARE @hasPermission BIT = 0;
-- Check user's clearance level
IF EXISTS (
SELECT 1 FROM user_permissions up
JOIN data_classification dc ON dc.classification_level = up.max_clearance_level
WHERE up.user_id = @userId
AND dc.table_name = @tableName
AND up.max_clearance_level >= @classificationLevel
)
BEGIN
SET @hasPermission = 1;
END
RETURN @hasPermission;
END;
-- View for secure data access
CREATE VIEW secure_user_data AS
SELECT
CASE
WHEN dc.masking_required = 1 THEN
CASE dc.column_name
WHEN 'ssn' THEN '***-**-' + RIGHT(ud.ssn, 4)
WHEN 'phone' THEN '+1-***-***-' + RIGHT(ud.phone, 4)
WHEN 'email' THEN LEFT(ud.email, 3) + '***@' +
SUBSTRING(ud.email, CHARINDEX('@', ud.email) + 1, LEN(ud.email))
ELSE ud.email
END
ELSE ud.email
END AS masked_email,
ud.user_id,
ud.created_date
FROM user_data ud
JOIN data_classification dc ON dc.table_name = 'user_data'
AND dc.column_name = 'email'
WHERE dbo.CheckDataAccessPermission(SYSTEM_USER, 'user_data', dc.classification_level) = 1;
5. Consent Management
Consent Tracking System:
-- Consent management table
CREATE TABLE user_consent (
consent_id BIGINT IDENTITY(1,1) PRIMARY KEY,
user_id VARCHAR(50),
consent_type VARCHAR(50), -- MARKETING, ANALYTICS, THIRD_PARTY_SHARING
consent_status VARCHAR(20), -- GRANTED, DENIED, WITHDRAWN
consent_date DATETIME,
withdrawal_date DATETIME,
consent_version VARCHAR(10),
ip_address VARCHAR(45),
user_agent NVARCHAR(500),
legal_basis VARCHAR(100), -- LEGITIMATE_INTEREST, CONSENT, CONTRACT
data_usage_purpose NVARCHAR(500)
);
-- Function to check active consent
CREATE FUNCTION HasActiveConsent(
@userId VARCHAR(50),
@consentType VARCHAR(50)
)
RETURNS BIT
AS
BEGIN
DECLARE @hasConsent BIT = 0;
IF EXISTS (
SELECT 1 FROM user_consent
WHERE user_id = @userId
AND consent_type = @consentType
AND consent_status = 'GRANTED'
AND withdrawal_date IS NULL
)
BEGIN
SET @hasConsent = 1;
END
RETURN @hasConsent;
END;
-- Stored procedure for data processing with consent check
CREATE PROCEDURE ProcessUserDataWithConsent
@userId VARCHAR(50),
@dataPurpose VARCHAR(100)
AS
BEGIN
SET NOCOUNT ON;
-- Check if user has consented to this type of data processing
IF dbo.HasActiveConsent(@userId, @dataPurpose) = 0
BEGIN
RAISERROR('User has not provided consent for %s data processing', 16, 1, @dataPurpose);
RETURN;
END
-- Log the data processing activity
INSERT INTO data_processing_log (
user_id, processing_purpose, processing_date, consent_verified
)
VALUES (@userId, @dataPurpose, GETDATE(), 1);
-- Proceed with data processing
-- ... actual processing logic here
END;
6. Data Quality and Integrity
Data Validation Framework:
-- Data quality rules table
CREATE TABLE data_quality_rules (
rule_id INT IDENTITY(1,1) PRIMARY KEY,
table_name VARCHAR(100),
column_name VARCHAR(100),
rule_type VARCHAR(50), -- FORMAT, RANGE, UNIQUENESS, COMPLETENESS
rule_definition NVARCHAR(MAX),
severity VARCHAR(20), -- ERROR, WARNING, INFO
is_active BIT DEFAULT 1
);
-- Function to validate email format
CREATE FUNCTION ValidateEmailFormat(@email VARCHAR(255))
RETURNS BIT
AS
BEGIN
DECLARE @isValid BIT = 0;
-- Basic email validation pattern
IF @email LIKE '%_@_%._%'
AND @email NOT LIKE '%@%@%'
AND @email NOT LIKE '%..%'
AND @email NOT LIKE '%.@%'
AND @email NOT LIKE '%@.%'
BEGIN
SET @isValid = 1;
END
RETURN @isValid;
END;
-- Stored procedure for data quality monitoring
CREATE PROCEDURE MonitorDataQuality
AS
BEGIN
SET NOCOUNT ON;
-- Check for data quality violations
INSERT INTO data_quality_violations (
rule_id, table_name, column_name, record_id,
violation_details, detected_date
)
SELECT
dqr.rule_id,
dqr.table_name,
dqr.column_name,
ud.user_id,
'Data quality rule violation: ' + dqr.rule_definition,
GETDATE()
FROM data_quality_rules dqr
CROSS APPLY (
SELECT user_id, email
FROM user_data
WHERE dqr.column_name = 'email'
AND dqr.rule_type = 'FORMAT'
AND dbo.ValidateEmailFormat(email) = 0
) ud
WHERE dqr.is_active = 1;
END;
7. Compliance Reporting
GDPR Compliance Reporting:
-- GDPR compliance reporting view
CREATE VIEW gdpr_compliance_report AS
SELECT
'Data Retention' AS compliance_area,
COUNT(*) AS record_count,
CASE
WHEN MAX(created_date) < DATEADD(YEAR, -1, GETDATE())
THEN 'REQUIRES_REVIEW'
ELSE 'COMPLIANT'
END AS status
FROM user_data
WHERE consent_status = 'GRANTED'
UNION ALL
SELECT
'Consent Management' AS compliance_area,
COUNT(*) AS record_count,
CASE
WHEN COUNT(CASE WHEN consent_status = 'GRANTED' THEN 1 END) > 0
THEN 'COMPLIANT'
ELSE 'NON_COMPLIANT'
END AS status
FROM user_consent
WHERE consent_date >= DATEADD(YEAR, -1, GETDATE())
UNION ALL
SELECT
'Data Access Audit' AS compliance_area,
COUNT(*) AS record_count,
CASE
WHEN COUNT(*) > 0
THEN 'COMPLIANT'
ELSE 'NO_AUDIT_DATA'
END AS status
FROM data_access_audit
WHERE access_timestamp >= DATEADD(MONTH, -6, GETDATE());
8. Data Breach Response
Incident Response Framework:
-- Data breach incident tracking
CREATE TABLE data_breach_incidents (
incident_id BIGINT IDENTITY(1,1) PRIMARY KEY,
incident_date DATETIME,
detection_date DATETIME,
breach_type VARCHAR(50), -- UNAUTHORIZED_ACCESS, DATA_EXPOSURE, SYSTEM_BREACH
affected_records_count INT,
data_types_affected NVARCHAR(500),
severity_level VARCHAR(20), -- LOW, MEDIUM, HIGH, CRITICAL
containment_status VARCHAR(20),
notification_sent BIT DEFAULT 0,
regulatory_reporting_required BIT DEFAULT 0
);
-- Stored procedure for breach response
CREATE PROCEDURE HandleDataBreach
@breachType VARCHAR(50),
@affectedTable VARCHAR(100),
@severityLevel VARCHAR(20)
AS
BEGIN
SET NOCOUNT ON;
DECLARE @affectedCount INT;
-- Assess impact
SET @affectedCount = (SELECT COUNT(*) FROM INFORMATION_SCHEMA.TABLES
WHERE TABLE_NAME = @affectedTable);
-- Log incident
INSERT INTO data_breach_incidents (
incident_date, detection_date, breach_type,
affected_records_count, severity_level
)
VALUES (GETDATE(), GETDATE(), @breachType, @affectedCount, @severityLevel);
-- Immediate containment actions
IF @severityLevel IN ('HIGH', 'CRITICAL')
BEGIN
-- Freeze affected accounts
UPDATE user_data
SET account_status = 'SUSPENDED'
WHERE user_id IN (
SELECT DISTINCT user_id
FROM data_access_audit
WHERE access_timestamp >= DATEADD(HOUR, -24, GETDATE())
);
-- Log containment action
INSERT INTO incident_response_log (
incident_id, action_taken, action_timestamp
)
VALUES (@@IDENTITY, 'Account suspension for affected users', GETDATE());
END
-- Trigger notification if required
IF @severityLevel IN ('MEDIUM', 'HIGH', 'CRITICAL')
BEGIN
EXEC SendBreachNotification @breachType, @severityLevel, @affectedCount;
END
END;
Key Technical Implementation Points:
- Encryption at Rest and in Transit: Implement column-level encryption for sensitive data
- Access Control: Role-based access control with principle of least privilege
- Audit Trails: Comprehensive logging of all data access and modifications
- Data Minimization: Only collect and retain necessary data
- Consent Management: Granular consent tracking with withdrawal capabilities
- Data Quality: Validation rules to ensure data integrity
- Incident Response: Automated detection and response procedures
- Compliance Monitoring: Regular reporting on regulatory compliance
This comprehensive approach demonstrates technical leadership in implementing ethical data handling practices while maintaining system performance and usability.