IAT 1 Question Bank
On this page
CMC702 IAT-1 — Exam-Ready Answer Bank
Each answer below is written for a 10-mark answer. State assumptions before applying an algorithm; in MapReduce, every intermediate record is explicitly a (key, value) pair.
Module I
1. Define Big Data. Explain the 5 Vs (Volume, Velocity, Variety, Veracity, and Value) with suitable real-world examples.
Definition. Big Data is data whose size, speed of generation, diversity, and uncertainty exceed the practical ability of conventional database systems to capture, store, process, and analyse economically. It is therefore both a data problem and a systems problem: distributed storage, parallel processing, governance, and analytics must work together.
| V | Meaning | Example | Engineering consequence |
|---|---|---|---|
| Volume | Amount of data | An e-commerce site stores years of clickstreams, orders, images, and logs measured in terabytes or petabytes. | Partitioned, distributed storage such as HDFS/object storage; compression and lifecycle policies. |
| Velocity | Rate at which data is produced, moved, and acted upon | Sensors emit readings every second; a payment gateway must score a transaction in milliseconds. | Streaming ingestion, buffering, windowing, and low-latency processing rather than only overnight batches. |
| Variety | Different structures and formats | Relational orders, JSON events, free-text reviews, images, audio, and GPS points. | Schema-on-read, data lakes, format-aware processing, and multiple storage models. |
| Veracity | Trustworthiness and quality | Duplicate customer IDs, missing sensor values, fraudulent clicks, or conflicting addresses. | Validation, provenance, deduplication, anomaly detection, confidence scores, and access controls. |
| Value | Useful business or scientific outcome extracted from data | A recommendation increases conversion; traffic analysis reduces travel time; predictive maintenance prevents downtime. | A measurable objective, suitable features/queries, privacy controls, and a cost-benefit test. |
The Vs interact. High-velocity, low-veracity sensor data has little value until it is cleaned and correlated with location and time. Likewise, storing high-volume data without a use case creates cost rather than value. A practical pipeline is:
sources -> ingest -> validate/quality checks -> distributed storage
-> batch/stream processing -> analytics/model -> decision/action -> feedback
Thus, Big Data is not simply “large data”; it is the combination of scale, speed, complexity, uncertainty, and the ability to turn data into decisions. Hadoop’s distributed storage and parallel processing are examples of infrastructure designed for these characteristics (Apache Hadoop modules).
2. Differentiate between Traditional Data Processing and the Big Data Business Approach based on architecture, scalability, processing, storage, and applications.
| Aspect | Traditional data processing | Big Data business approach |
|---|---|---|
| Architecture | Centralised or scale-up server; a database is the principal system of record. | Distributed cluster/cloud: storage and compute are spread over many commodity or cloud nodes. |
| Scalability | Mainly vertical: add CPU, RAM, or storage to one machine; expensive at high scale. | Horizontal: add nodes and partition data; capacity grows incrementally. |
| Processing | SQL transactions and scheduled reports; strong schema and predictable workloads. | Parallel batch and stream jobs, SQL-on-data-lake, machine learning, and event-driven decisions. |
| Storage | Normalised tables, fixed schema, structured records; limited tolerance for very large unstructured files. | Data lake/distributed file system stores structured, semi-structured, and unstructured data; schema can be applied at read time. |
| Consistency/latency | ACID transactions and low-latency point updates are primary. | Often combines eventual consistency and high throughput with separate transactional systems where required. |
| Data arrival | Periodic ETL from known operational systems. | Continuous ingestion from applications, sensors, social media, devices, logs, and external feeds. |
| Business use | Payroll, inventory transaction, accounting, standard MIS reports. | Fraud detection, recommendations, predictive maintenance, real-time pricing, sentiment, and exploratory discovery. |
| Economics | Proprietary high-end hardware/software and capacity planned in advance. | Commodity/cloud infrastructure, open-source ecosystem, elastic capacity, and pay-for-use options. |
In the traditional approach, the business first designs a schema and then loads approved data. In the Big Data approach, it often retains raw data with metadata, then creates curated views for different questions. For example, a retailer can keep clickstream JSON, order tables, and review text in a lake; a batch job builds customer segments while a stream job flags abandoned carts.
The approaches are complementary, not mutually exclusive. A bank should retain an ACID core-banking database for balances, while using distributed analytics for fraud and risk. A sensible architecture is therefore polyglot: use each engine for the workload it handles safely. Hadoop provides distributed storage and computation, but it does not replace every OLTP database.
3. Explain the types of Big Data (Structured, Semi-structured, and Unstructured) with suitable examples from industry.
1. Structured data. Data has a fixed, explicitly defined schema: rows, columns, types, keys, and relationships. Examples are an ATM transaction table (account_id, time, amount, merchant_id), an airline reservation table, or an inventory table. It is easy to validate and query using SQL. The limitation is that a rigid schema is inconvenient when fields change or when the data is not naturally tabular.
2. Semi-structured data. Data has tags, keys, metadata, or nested structure but does not require every record to have identical columns. Examples include JSON product events, XML invoices, email headers, web-server logs, and IoT messages such as:
{"device":"D17", "time":"2026-08-12T10:00:00Z", "temp":31.4,
"location":{"lat":18.52,"lon":73.86}}
Two device types may add different fields while retaining common identifiers. A document store or schema-on-read lake is suitable; validation is still needed at ingestion.
3. Unstructured data. No predefined tabular model captures the content directly. Examples are CCTV video, medical X-rays, MRI images, call recordings, PDF documents, social-media text, photographs, and audio. Analysis requires specialised representations: OCR/text extraction, embeddings, image features, speech recognition, or video metadata. The original binary is usually stored in a file/object store; derived metadata can be placed in a database.
| Property | Structured | Semi-structured | Unstructured |
|---|---|---|---|
| Schema | Fixed and explicit | Flexible, keys/tags/nesting | Implicit or absent |
| Typical format | Tables, CSV | JSON, XML, logs | Images, video, audio, free text |
| Example | Bank ledger | Click event | CCTV stream |
| Typical processing | SQL/relational algebra | JSON/document queries, parsing | ML, search, OCR, media analytics |
A real retailer combines all three: structured orders, semi-structured click events, and unstructured review text/images. Big Data systems preserve their source format and attach common metadata such as customer, time, and provenance so that the forms can be joined during analysis.
4. Describe the Hadoop Architecture and explain the functions of its core components with a neat diagram.
Apache Hadoop is an ecosystem for distributed storage and parallel processing. Its current core modules are Hadoop Common, HDFS, YARN, and MapReduce (official module description). A simplified Hadoop 2/3 architecture is:
Client / applications
|
+--------------+--------------+
| |
HDFS client (metadata) YARN client submits job
| |
+--------v--------+ +-------v--------+
| Active NameNode | | ResourceManager|
| namespace/meta | | cluster scheduler|
+--------+--------+ +-------+--------+
| heartbeats/ | launches
| block reports | containers
+---------+---------+ +---------+---------+
| | | |
+-----v-----+ +-----v-----+ +--v----------+ +-----v------+
| DataNode | ... | DataNode | |NodeManager | |NodeManager |
| HDFS blocks| | HDFS blocks| |Application | |Application |
+-----------+ +-----------+ |Master/tasks| |Master/tasks |
+-------------+ +-------------+
- HDFS splits large files into blocks, replicates them over DataNodes, and gives clients high-throughput streaming access. The NameNode stores namespace and block-location metadata; DataNodes store and serve the blocks. Client data normally moves directly between client and DataNodes (HDFS design).
- YARN separates resource management from applications. The ResourceManager allocates cluster resources, NodeManagers monitor each worker, and an ApplicationMaster coordinates one submitted application.
- MapReduce is a batch execution model: mappers emit intermediate key-value pairs, the framework partitions/shuffles/sorts them, and reducers aggregate them. It runs in YARN containers and exploits data locality.
- Hadoop Common supplies shared libraries, configuration, RPC, serialization, and filesystem APIs.
For high availability, a deployment uses active/standby NameNodes; for metadata scale it may use federation. In older Hadoop 1 architecture, a JobTracker scheduled MapReduce and TaskTrackers ran tasks; YARN replaced that central job/resource role. A complete answer should distinguish storage (HDFS), resource scheduling (YARN), and computation (MapReduce) rather than calling all of them “HDFS.”
5. Explain the Hadoop Ecosystem and discuss the role of HDFS, MapReduce, Hive, Pig, HBase, Sqoop, and Flume in Big Data processing.
RDBMS --Sqoop--> HDFS <---Flume--- logs/events
|
+------------+----------------+
| |
MapReduce Hive / Pig
| |
+----------> HBase <----------+
|
reports / applications
| Component | Role and appropriate use |
|---|---|
| HDFS | Distributed, replicated file system for large files and high-throughput scans. Data is divided into blocks and stored on DataNodes under NameNode metadata. |
| MapReduce | Batch parallel processing over HDFS: map, shuffle/sort, and reduce. Suitable for large scans, joins, aggregation, and jobs where latency of minutes is acceptable. |
| Hive | Data-warehouse layer providing SQL-like HiveQL over large datasets; useful for ETL, reporting, and ad-hoc analytics, not low-latency OLTP (Apache Hive introduction). |
| Pig | High-level Pig Latin data-flow language for loading, filtering, grouping, joining, and transforming irregular data; the script is compiled into distributed jobs (Apache Pig). |
| HBase | Distributed column-family NoSQL database on Hadoop/HDFS for random, near-real-time reads and writes of sparse, very large tables; it is not a replacement for Hive’s analytical scans (Apache HBase overview). |
| Sqoop | Bulk, parallel import/export between relational databases and Hadoop. Use it for structured historical tables, not high-rate event streams (Sqoop guide). |
| Flume | Distributed event/log collector. Its normal flow is Source → Channel → Sink, with buffering and fan-in from many log sources (Flume User Guide). |
A typical pipeline is: Flume collects web logs into HDFS; Sqoop imports customer/order tables; Hive or Pig cleans and joins them; MapReduce executes the heavy batch plan; HBase serves a selected operational result, such as a customer profile or alert. YARN supplies resource management, although it is not listed in the question. Components should be chosen by workload: HDFS for durable bulk storage, HBase for keyed serving, and Hive/MapReduce for scans and analytics.
6. Analyze a Big Data case study in the healthcare, banking, or e-commerce domain. Explain the challenges faced and how Hadoop-based solutions addressed them.
Case: hospital readmission and patient-risk analytics. A hospital network wants to predict 30-day readmission and monitor chronic patients. Requirements are to combine years of electronic health records (EHR), laboratory results, prescriptions, doctor notes, claims, wearable streams, and medical images; preserve auditability; protect personally identifiable information (PII); and produce both daily population reports and near-real-time alerts.
Challenges. EHR tables are structured but use inconsistent codes; notes are free text; devices produce high-velocity data; images are very large; records may be incomplete or duplicated; access is highly sensitive; and patient identity must be resolved without exposing raw identifiers. Conventional single-server ETL becomes slow and expensive for repeated historical scans.
Hadoop-based design and flow:
EHR/claims DB --Sqoop--+ +--> Hive/MapReduce
v | clean/join/features
wearables/logs --Flume-> HDFS raw -> curated --> +--> ML/risk scores
notes/images ----------> (encrypted, ACLs) +--> HBase alerts/profile
- Sqoop imports partitioned historical tables incrementally; Flume collects device/application events. Raw files are immutable and timestamped.
- HDFS replicates blocks and stores text, JSON, and large binary/object references. A de-identification job replaces direct identifiers with governed tokens; HDFS permissions, encryption, audit logs, and retention policies enforce least privilege.
- MapReduce/Hive parses codes, removes duplicates, normalises units, extracts note terms, and joins records by the tokenised patient ID and time. Missing values and data-quality flags are retained rather than silently fabricated.
- Features such as prior admissions, medication gaps, abnormal lab trends, and recent sensor statistics are computed in batch. Risk scores are written to HBase keyed by
(patient_token, event_time)for fast clinical lookup; a clinician-facing service applies an alert threshold. - Governance prevents the model from becoming an unaudited diagnosis: validate against held-out patients, monitor bias and drift, record feature provenance, and require clinician review.
Benefits/trade-offs. Distributed storage makes repeated population scans feasible and keeps raw plus derived data together. Parallel jobs handle heterogeneous historical data at lower scale-out cost. However, classic MapReduce is not a millisecond clinical decision engine; streaming/serving components are needed, and replication increases storage cost. Hadoop assists computation but does not itself solve medical privacy, identity, data semantics, or clinical validation.
7. Evaluate the advantages and limitations of adopting Hadoop for processing large-scale enterprise data. Support your answer with suitable examples.
Advantages
- Horizontal scalability: files and work are partitioned across nodes; a cluster can grow by adding machines.
- Fault tolerance: HDFS replication and task re-execution tolerate failed disks, workers, and individual tasks. DataNodes send heartbeats and block reports; under-replicated blocks can be recreated (HDFS design).
- Cost and openness: commodity/cloud hardware and open-source APIs reduce dependence on a proprietary appliance.
- Data variety: HDFS can retain large structured, semi-structured, and unstructured source files for later schema-on-read analysis.
- Data locality and throughput: computation can be scheduled near blocks, reducing network movement for full scans.
- Ecosystem: Hive, HBase, Flume, Sqoop, and other tools support ingestion, analytics, and serving.
- Reproducibility: immutable raw data plus repeatable jobs helps reprocess data when business rules change.
Limitations
- HDFS and MapReduce favour large sequential batch jobs; small files, random updates, and interactive millisecond queries are poor fits.
- Replication consumes storage and network bandwidth; cluster operations, security, upgrades, and capacity planning require specialists.
- MapReduce writes intermediate results to disk and has job startup/shuffle overhead, so iterative ML and interactive workloads are slow compared with Spark or specialised engines.
- Classic Hadoop is not an OLTP database: row-level ACID transactions, constraints, and rich joins are not its primary strength.
- Data governance is difficult when raw data is copied into many derived datasets; schema, lineage, PII deletion, and access policy must be engineered.
- A cluster can be underutilised for irregular or small workloads; cloud object storage and managed query engines may be cheaper.
Evaluation. Use Hadoop for a retailer’s multi-year clickstream aggregation, genome batch analysis, or log archive; use an RDBMS for account balances, HBase/document stores for keyed serving, and Spark/streaming systems for low-latency iterative analysis. Hadoop is therefore a durable scale-out foundation, not automatically the best answer for every “big” dataset. The decision should compare data size/rate, latency, query shape, consistency, team skill, total cost, and governance.
8. Compare Big Data technologies with traditional database systems. Justify why organizations are migrating towards Hadoop-based solutions.
| Dimension | Traditional relational database | Hadoop-oriented Big Data platform |
|---|---|---|
| Data model | Normalised tables, fixed schema, relations and constraints | Files/tables/documents/streams; schema-on-read and multiple formats |
| Main strength | ACID transactions, joins, indexes, low-latency point queries | High-throughput distributed storage and parallel scans |
| Scaling | Usually scale up; sharding is possible but operationally specialised | Scale out by partitioning blocks and jobs across nodes |
| Data | Mostly structured operational data | Structured plus logs, JSON, text, images, and other large files |
| Processing | SQL and transaction/report workloads | MapReduce/Hive/Pig and other batch/parallel frameworks; Hadoop itself is primarily batch-oriented |
| Failure model | Hardware redundancy often hidden in appliance/cluster | Failure is expected; replication and task re-execution are explicit design features |
| Schema/change | Schema-first; migration can be expensive | Raw retention and schema-on-read allow exploratory use, but governance remains essential |
| Best examples | Banking ledger, order commit, inventory update | Clickstream history, log analytics, feature generation, archive, large joins |
Organizations migrate or augment systems because data sources have multiplied, traditional vertical scaling is expensive, historical retention is growing, and analysts need to combine data that was formerly discarded. Hadoop can store raw data cheaply, process it in parallel, and support exploratory analytics before a final schema is known. A retailer might load all click events and reviews into HDFS, compute segments with Hive/MapReduce, and publish only a compact customer table to the relational system.
Migration is not a blanket replacement. The justified target architecture is usually hybrid: relational databases continue to own transactions; Hadoop handles large-scale ingestion, batch analytics, and data-lake retention; HBase or another serving store handles high-volume keyed access. Migration should include data quality, encryption, lineage, backup, skills, workload benchmarks, and a rollback plan. Moving simply because “Hadoop is big” can increase complexity without business value.
9. Case Study: A smart city generates huge amounts of traffic, CCTV, GPS, and weather data every minute. Design a Hadoop-based architecture for storing and processing this data. Justify the selection of Hadoop components.
Requirements. Ingest continuous heterogeneous data; retain raw data for investigations; compute congestion and incident summaries every few minutes; process historical trends; support spatial/time queries; tolerate node failures; and restrict CCTV/PII access. A raw video frame is not forced into a row: store video in a suitable object/file format and keep indexed metadata in Hadoop-compatible stores.
traffic loops/GPS/weather CCTV gateways
| |
Flume/Kafka* Flume/collector
+-----------> HDFS raw landing zone
(date/hour/source partitions)
|
+--------------------+-------------------+
| |
Hive/MapReduce batch stream/short jobs*
clean, join, aggregate minute congestion
| |
curated Parquet/Hive tables HBase serving: (road,time)
| |
reports, GIS/ML, planning dashboard/alerts/API
* A modern deployment may use Kafka and Spark/Flink for streaming; if the syllabus requires Hadoop ecosystem components, Flume + MapReduce/Hive are the corresponding choices.
Components and data flow. Flume collectors buffer logs and events; Sqoop periodically imports structured road-inventory or weather history from an RDBMS. HDFS stores replicated blocks partitioned by event date, city zone, and source. MapReduce cleans timestamps, deduplicates GPS points, joins readings to road segments, and computes vehicle count, mean speed, queue length, and weather correlation. Hive exposes SQL-like tables for planners. HBase stores the latest (road_id, time_bucket) status for fast dashboard lookup; a spatial index or derived road-segment key avoids scanning all CCTV data. Video is access-controlled and compressed; object/file paths and recognised-event metadata are analysed rather than exposing all footage.
Justification/trade-offs. HDFS is suitable for large sequential CCTV/history scans and replication protects against node loss. MapReduce/Hive provide economical batch trends; HBase provides low-latency keyed status. Partitioning by time and zone enables pruning, while compression reduces storage. Exact real-time alerts need a streaming engine, not a multi-minute MapReduce job. Governance must cover camera access, retention, encryption, anonymisation, and audit. Thus Hadoop is the historical and batch backbone, complemented by streaming and serving layers.
10. Application-Based Question: An online shopping company wants to analyze customer purchasing behaviour during festive sales. Propose a Big Data solution using the Hadoop ecosystem. Draw a suitable architecture and explain the data flow from data collection to analytics.
Requirements. Capture clicks, searches, carts, orders, payments, reviews, campaign exposure, and product catalogue data; handle a sales spike; compute funnels, cohorts, basket associations, customer value, and recommendations; and preserve privacy while keeping operational checkout reliable.
web/app events, logs ----Flume----+
orders/catalog RDBMS -----Sqoop----+--> HDFS raw (immutable, partitioned)
reviews/JSON ----------------------+ |
v
Hive/MapReduce/Pig ETL
parse -> clean -> identity/time join
|
+-------------------------------+----------------+
| | |
Hive curated facts MapReduce jobs HBase serving
sales/funnel/cohorts segments/associations customer/product view
| | |
BI dashboards recommendations campaign API
Flow. Flume receives clickstream events through source-channel-sink agents and buffers bursts. Sqoop performs incremental imports of orders and catalogue tables using a stable key/time boundary. HDFS keeps raw events partitioned by date, region, and event type; access controls tokenise customer identity. Hive external tables define schemas over curated files. Pig or MapReduce normalises events, removes duplicate retries, orders events by customer/session, and joins orders with products and campaign exposures.
Analytics include: (i) conversion funnel view → cart → checkout → paid, (ii) cohort retention, (iii) average order value and customer lifetime features, (iv) association counts such as products bought together, and (v) recommendation candidates. A job can emit (customer_id, feature_vector) and write selected current features to HBase for fast campaign lookup. Hive stores aggregate facts for analysts; MapReduce handles large scans; HBase handles keyed online retrieval.
During the sale, scale collectors and use partitioned files; do not put checkout commits in a batch Hadoop job. Enforce consent, masking, retention, and role-based access for payment/customer fields. The trade-off is high throughput and historical insight versus batch latency and ecosystem complexity; a real deployment may add Kafka and Spark for streaming but keeps the Hadoop data-lake pattern.
Module II
1. Explain the architecture of Hadoop Distributed File System (HDFS) with a neat diagram. Describe the functions of NameNode, DataNode, and Secondary NameNode.
HDFS client
1. metadata request
|
+---------v----------+
| NameNode |
| namespace, fsimage |
| edits, block map |
+----+----------+----+
| |
2. locations/commands | checkpoints
| v
+----------+--+ +---+----------------+
| DataNode A | | Secondary NameNode |
| blocks 1,4 | | fetch fsimage+edits|
+------+------+ | merge checkpoint |
| +--------------------+
+------v------+ (not hot standby)
| DataNode B | ... DataNode C
| replica blocks replica blocks
+-------------+
A large file is divided into configurable blocks. The NameNode is the master: it maintains the filesystem namespace (directories, names, permissions), file-to-block mapping, and block-to-DataNode locations in memory and persistent metadata. It does not normally carry user file bytes. It accepts client metadata operations, chooses replica placement, receives DataNode heartbeats/block reports, and initiates re-replication or deletion.
A DataNode is a worker/storage daemon. It stores block replicas on local disks, serves client read/write requests, verifies checksums, participates in the write pipeline, and reports health and block inventory. For a read, the client asks the NameNode for locations and then reads directly from a suitable DataNode. For a write, the NameNode selects targets and DataNodes pipeline replicas.
The traditional Secondary NameNode is a checkpoint helper, not a backup/automatic failover NameNode. It periodically obtains the current fsimage and edit log, merges them, and returns a new checkpoint, limiting edit-log growth and restart time. Modern HA uses active/standby NameNodes and shared or journaled edit storage instead (HDFS design; HDFS HA). HDFS is optimised for high-throughput streaming and large files, not many tiny random updates.
2. Explain the physical organization of compute nodes in HDFS. How does HDFS achieve reliability and fault tolerance?
A cluster is physically organised as racks containing worker machines. A DataNode generally has multiple local disks; it stores HDFS blocks and may also run YARN containers. A client or compute task is preferably placed on the same node as the needed block, then in the same rack, then in another rack. This is data locality: it reduces network traffic while preserving a rack-failure margin.
Rack 1 Rack 2
+------------------+ +------------------+
| DN1: B1, B4 | | DN3: B1, B2 |
| DN2: B2, B5 | | DN4: B3, B4 |
+------------------+ +------------------+
NameNode metadata: B1->{DN1,DN3}, B2->{DN2,DN3}, ...
Reliability mechanisms
- Block splitting: a file is split into independent blocks, so one disk failure affects only those blocks rather than an entire monolithic file.
- Replication: each block has a configurable replication factor, commonly 3 in examples. Replicas are placed across nodes and preferably racks. If one copy fails, another serves reads.
- Checksums: DataNodes detect corrupt blocks; a client can read another verified replica.
- Heartbeats and block reports: DataNodes periodically report liveness and stored blocks to the NameNode. Missing heartbeats mark a node dead; the NameNode schedules new replicas for its blocks (HDFS design).
- Write pipeline and acknowledgements: a client writes to a chain of DataNodes; each forwards the packet and acknowledgements return through the chain. A failed target is removed and the pipeline is rebuilt.
- Task re-execution/data locality: YARN/MapReduce can rerun a failed task, often on another node, because input remains in HDFS.
- Metadata protection: checkpoints, edit logs, active/standby NameNodes, and journal/quorum mechanisms reduce NameNode failure risk.
Replication protects availability, not logical mistakes or ransomware. Therefore use snapshots, backups, access control, encryption, monitoring, and tested recovery. Replication factor 3 gives node/rack resilience but costs roughly three times raw block capacity before compression.
3. Describe the MapReduce programming model. Explain the sequence of Map, Shuffle & Sort, and Reduce phases with a suitable example.
MapReduce transforms input key-value records into output key-value records:
InputFormat: <K1,V1>
|
v
Mapper: <K1,V1> -> zero or more <K2,V2>
|
v
Partition + shuffle + sort/group: all same K2 together
|
v
Reducer: <K2, Iterable<V2>> -> zero or more <K3,V3>
|
v
OutputFormat: final files in HDFS
Example: word count. Suppose lines are big data and big systems. TextInputFormat gives approximately (offset, line) records. Mappers emit:
("big",1), ("data",1), ("big",1), ("systems",1)
The partitioner sends a key to one reducer, ensuring all occurrences of big meet. Shuffle copies each mapper’s partition to the assigned reducer. Sort orders intermediate keys, and grouping produces:
"big" -> [1,1]
"data" -> [1]
"systems" -> [1]
The reducer computes sum(values) and emits ("big",2), ("data",1), ("systems",1). Hadoop’s reducer receives one key and an iterable of all values for that key; the shuffle fetches mapper partitions and groups keys (MapReduce tutorial).
Execution details. Input files are split; normally one map task processes each split. Map output is buffered, partitioned, sorted, and spilled. Reducers fetch their partitions, merge-sort them, then call reduce once per grouped key. Number of reducers controls parallelism and output files. A combiner may perform local aggregation before transfer, but it is optional. The model is excellent for associative/commutative aggregations and embarrassingly parallel transformations; joins and iterative algorithms may require multiple jobs.
4. Explain the role of Mapper, Reducer, and Combiner in MapReduce. Illustrate with a Word Count example.
Mapper. Converts each input record (K1,V1) into zero or more intermediate (K2,V2) pairs. For word count, the key is the normalised word and the value is integer 1:
map(offset, line):
for word in tokenize(lowercase(line)):
emit(word, 1)
Reducer. Receives one grouped key and all values assigned to it, then emits a final result:
reduce(word, counts):
total = 0
for c in counts: total = total + c
emit(word, total)
Combiner. An optional local mini-reducer runs after a mapper, before shuffle, to reduce network traffic. For a mapper that saw big big data, it can turn (big,1),(big,1),(data,1) into (big,2),(data,1). The reducer still adds values arriving from all mappers. It may execute zero, one, or multiple times, so the operation must be associative/commutative and must not rely on side effects; it is not a replacement for the reducer (Hadoop MapReduce tutorial).
Mapper A: (big,1),(big,1),(data,1) --combiner--> (big,2),(data,1)
Mapper B: (big,1),(systems,1) --combiner--> (big,1),(systems,1)
shuffle/group
Reducer: big->[2,1] => (big,3)
data->[1] => (data,1)
systems->[1]=> (systems,1)
Key-value type logic: with text input, mapper input may be (byte_offset, Text line); mapper output is (Text word, IntWritable 1). Combiner input/output is (Text, IntWritable) for this sum. Reducer input is (Text, Iterable<IntWritable>); final output is (Text, IntWritable). The partitioner hashes the word, so equal words cannot be split across reducers. A combiner improves performance, not correctness: if it is omitted, the reducer receives the original ones and obtains the same totals.
5. Explain how Hadoop handles node failures during MapReduce execution. Why is fault tolerance an important feature of Hadoop?
Hadoop assumes that failures of disks, nodes, processes, and network paths are normal at cluster scale. YARN monitors containers and NodeManagers; MapReduce monitors task attempts and job progress.
Failure handling sequence
- A mapper or reducer process crashes, times out, or stops reporting progress. The framework marks that task attempt failed.
- It launches a replacement attempt, preferably where input blocks are local. A failed map can reread its HDFS input; a failed reduce refetches the relevant map partitions.
- Intermediate map output is temporary. If its node disappears, the map task is rerun and its partition regenerated.
- If a reducer fails, other reducers continue; its assigned partitions are fetched again after the replacement starts. Output is committed atomically/through task-attempt management so speculative or failed attempts do not become the final result.
- Speculative execution may run a duplicate attempt for a straggler. The first successful attempt wins, useful when one slow disk/node delays an otherwise complete job, though it consumes resources and should be disabled for side-effecting external work.
- HDFS detects dead DataNodes via missed heartbeats and re-replicates under-replicated blocks. Checksums let clients avoid corrupt replicas.
- Configurable retry limits, health checks, and skipped-record options handle repeatedly bad records; counters and logs support diagnosis.
Fault tolerance matters because a job using thousands of machines has a high probability that something fails during its runtime. Without re-execution and replicated input, a single worker failure would lose the entire multi-hour computation. The trade-off is extra storage, network traffic, and possibly duplicate work. Exactly-once effects are not automatic for external side effects, so tasks should be deterministic and write through Hadoop’s controlled output mechanisms. This is why Hadoop can provide reliable batch results on commodity hardware rather than requiring every machine to be perfect (MapReduce tutorial).
6. Explain the MapReduce algorithm for Matrix–Vector Multiplication. Illustrate the Map and Reduce functions with an example.
Let M be an m × n matrix and v an n-element column vector. We want y = Mv, where
Assumptions and representation. The matrix is supplied as sparse records (i, j, aij) and the vector is small enough to be distributed to every mapper (for example through a distributed cache). The key is the row index i; the value is a partial product.
map((i,j,aij), vector v):
emit(i, aij * v[j])
reduce(i, partial_products):
y_i = sum(partial_products)
emit(i, y_i)
The map phase can process each nonzero matrix entry independently. The shuffle groups all products for row i at one reducer. The reducer adds them. If the vector is too large for cache, partition/broadcast it through an additional join job; the core logic remains keying products by row.
Worked example
Matrix records and mapper output:
(1,1,1) -> (1, 1*5=5) (1,2,2) -> (1, 2*6=12)
(2,1,3) -> (2, 3*5=15) (2,2,4) -> (2, 4*6=24)
After grouping: 1 -> [5,12], 2 -> [15,24]. Reducers emit (1,17) and (2,39), hence y=[17,39]^T.
For a dense matrix this emits mn products and may be expensive; sparse representation emits only nonzeros. A combiner using addition is safe and can sum partial products locally. The algorithm parallelises rows and tolerates task failure because matrix blocks remain in HDFS. This is the standard MapReduce pattern described in Mining of Massive Datasets, Chapter 2 (book PDF).
7. Explain how the relational algebra operations Selection, Projection, Union, Intersection, and Difference are implemented using MapReduce.
Assume set semantics unless stated otherwise, and represent every tuple as a serialisable record. A mapper emits (key, value); the key is chosen so the reducer sees all records that must be compared.
| Operation | Map logic | Reduce logic |
|---|---|---|
Selection σ_p(R) | For each input tuple r, if predicate p(r) is true, emit (constant, r) or directly write r. | No reducer is necessary if the predicate is local. If one output file is required, a trivial reducer collects/outputs values. Example: filter amount>1000. |
Projection π_A(R) | Compute projected tuple t=r[A]; emit (t, null). | Output key t once and discard repeated values. This removes duplicates under set projection. |
Union R ∪ S | For every tuple t from either relation, emit (t, relation_tag); use the full tuple as key. | Output t once if it appears in at least one input. For bag union, retain multiplicities instead. |
Intersection R ∩ S | Emit (t, R) for R and (t, S) for S. | Output t only if the set of tags contains both R and S. |
Difference R − S | Emit (t, R) for R and (t, S) for S. | Output t only if tag R appears and tag S does not. |
Example: R={(1,a),(2,b)} and S={(2,b),(3,c)}. For union, key (2,b) receives [R,S] and is output once, giving {(1,a),(2,b),(3,c)}. For intersection it is the only key with both tags, giving {(2,b)}. For R−S, (1,a) has only R and is output; (2,b) is suppressed.
Selection can be map-only because each tuple is tested independently. Projection and set operations require a shuffle so equal tuple keys meet. A custom partitioner may distribute tuples; the invariant is that identical tuple keys go to the same reducer. If inputs are bags, union counts occurrences and intersection/difference need multiplicity rules rather than simple tag presence. MapReduce implementations of these operations are treated in Mining of Massive Datasets, Chapter 2 (Stanford book).
8. Evaluate the limitations of Hadoop in processing modern Big Data applications. Suggest suitable improvements or alternative technologies.
Limitations
- Latency: MapReduce job startup, disk spills, shuffle, and HDFS writes make it unsuitable for millisecond APIs or interactive exploration.
- Iterative computation: Machine learning and graph algorithms repeatedly read/write the same data; MapReduce’s materialisation at every iteration is expensive.
- Small files: Many tiny files overload NameNode metadata and create one-task-per-split overhead; HDFS favours large files.
- Random updates/OLTP: HDFS is write-once/append-oriented and is poor for frequent row updates, constraints, and multi-row ACID transactions.
- Operational complexity: Kerberos, permissions, HA, upgrades, capacity, data locality, and multiple ecosystem versions require specialist administration.
- Resource and data movement cost: replication, cross-rack shuffle, and serialisation consume disk/network; skewed keys create stragglers.
- Governance: a data lake can become a “data swamp” without schema, catalogue, lineage, quality, retention, and PII deletion controls.
- Technology fit: unbounded event processing, graph traversals, vector search, and low-latency serving need specialised engines.
Improvements/alternatives
- Use Spark for in-memory DAGs and iterative analytics, while retaining HDFS/object storage; use Tez or an optimised Hive execution engine to reduce unnecessary stages.
- Use Kafka + Flink/Spark Structured Streaming for event-time windows, state, and low-latency streaming.
- Use HBase/Cassandra/document stores for partition-keyed random access; use an RDBMS for transactional truth.
- Compact small files into columnar Parquet/ORC, partition by query predicates, compress, and use a metastore/catalogue.
- Reduce skew with better partition keys, combiners, salting hot keys, and approximate algorithms where exactness is unnecessary.
- Use managed cloud object storage and autoscaling where operations dominate, with encryption and governance built in.
Hadoop remains suitable for durable bulk storage and large batch scans, but modern architecture is often a lakehouse/polyglot platform. Select technology from requirements—latency, consistency, throughput, query shape, cost, and compliance—not from data size alone.
9. Case Study: A telecom company receives billions of call records every day. Design a MapReduce-based solution to calculate the total call duration for each customer. Explain the Map and Reduce logic.
Input assumption. Each CDR is a valid CSV/JSON record such as (call_id, caller_id, receiver_id, start_time, duration_seconds, status). Define “customer total” explicitly. Here, total duration is attributed to the billed caller caller_id, and only completed, non-duplicate calls with duration_seconds ≥ 0 are counted. If the requirement is both parties’ usage, emit one pair for each party or maintain separate caller/receiver metrics.
HDFS CDR files
|
InputFormat splits records
|
Mapper: (offset, cdr) -> (caller_id, duration_seconds)
|
optional Combiner: (caller_id, [durations]) -> (caller_id, local_sum)
|
partition/hash caller_id -> shuffle/group
|
Reducer: (caller_id, all durations) -> (caller_id, total_seconds)
|
HDFS totals / Hive report
Pseudocode
map(_, cdr):
if cdr.status == "COMPLETED" and cdr.duration_seconds >= 0:
emit(cdr.caller_id, cdr.duration_seconds)
combine(customer, durations):
emit(customer, sum(durations))
reduce(customer, partial_durations):
total = 0
for d in partial_durations: total += d
emit(customer, total)
The key is caller_id, so all records for a customer reach one reducer; the value is duration. The combiner is safe because integer addition is associative and commutative, and it can greatly reduce shuffle bytes. The partitioner must consistently route the same customer key to one reducer.
Worked example. Records for A have durations 120, 60, and 300 seconds; B has 40 and 20. Mapper pairs are (A,120),(A,60),(A,300),(B,40),(B,20). A combiner on one mapper may emit (A,180) and another (A,300). Reducer A outputs (A,480) seconds = 8 minutes; B outputs (B,60) seconds = 1 minute. Invalid/duplicate records are removed before aggregation, preferably using a deterministic call_id deduplication job if retries can create duplicates. Partition input by day and retain per-day outputs if auditing or incremental recomputation is needed.
10. Application-Based Question: An e-commerce company stores purchase records in HDFS. Design a MapReduce solution to identify the top 10 best-selling products. Draw the processing flow and justify your approach.
Assumption. “Best-selling” means the largest number of successfully paid units in a stated time interval, not revenue. Each purchase line is (order_id, product_id, quantity, status, timestamp); cancelled/refunded lines are excluded. If revenue is intended, emit unit_price × quantity and use a revenue metric instead.
HDFS purchase lines
|
Mapper 1: filter valid -> (product_id, quantity)
|
Combiner 1: local sum -> (product_id, local_units)
|
Reducer 1: sum all -> (product_id, total_units)
|
Mapper 2: maintain local 10 largest -> (constant, (product, units))
|
Reducer 2: merge candidates, sort -> global top 10
|
HDFS/report/API
Job 1 pseudocode
map(_, line):
p = parse(line)
if p.status == "PAID" and p.quantity > 0:
emit(p.product_id, p.quantity)
combine(product, quantities):
emit(product, sum(quantities))
reduce(product, quantities):
emit(product, sum(quantities))
This creates the exact total (product_id, units); the product key guarantees all its lines meet at one reducer. Job 2 avoids sending every product to one global reducer prematurely. Each mapper reads totals and keeps a min-heap of its ten largest products, emitting those ten as (TOP, (product_id, units)). A single reducer merges the at-most 10 × number_of_mappers candidates, sorts descending by units, applies a deterministic product-ID tie-breaker, and outputs ten.
The local-top-10 optimisation is correct: a product not in its mapper’s local top ten cannot be in the global top ten when that mapper’s records are partitioned after Job 1; every product has exactly one total record. If Job 1 already distributes totals, Job 2’s mapper may see arbitrary subsets, but the same local-candidate argument holds for the union of subsets. For skew or enormous output, use multiple final reducers with range partitioning and a tiny final merge. Keep a date/region partition and output totals plus job version for auditability.
Module III
1. Define NoSQL. Explain the need for NoSQL databases and discuss the major business drivers for adopting NoSQL in Big Data applications.
Definition. NoSQL (“not only SQL”) is a family of non-relational database systems designed around models such as key-value, document, wide-column, and graph. NoSQL does not mean “no query language” or “no consistency”; systems make different trade-offs in schema flexibility, distribution, transactions, and query patterns.
Why needed. Relational normalisation and joins are excellent for structured transactions, but a globally distributed application may need very high write throughput, horizontal scale, flexible/nested records, predictable access by a partition key, or a graph traversal. Forcing all data into one rigid relational schema can cause expensive migrations, join-heavy reads, and a single scaling bottleneck.
Business drivers
- Web/mobile scale: millions of users and bursty read/write traffic require scale-out and elastic capacity.
- Availability and global reach: replicas near users keep a service operating despite node/region failure; some workloads accept eventual consistency.
- Agile schema evolution: catalogues, profiles, and event payloads add fields frequently; documents can evolve without a table migration.
- Big and fast data: logs, IoT, social events, and clickstreams need sustained ingestion and partitioned storage.
- Application-oriented performance: a key-value lookup, document aggregate, column-family time series, or graph neighbourhood can avoid costly joins.
- Developer productivity and economics: JSON APIs map naturally to documents, while commodity/cloud nodes provide incremental scale.
- Specialised structure: relationships, sparse attributes, time-series partitions, and recommendation graphs each have different access patterns.
NoSQL is chosen by workload, not fashion. The design must specify partition key, query paths, consistency, retention, backup, and failure behaviour. For example, a social timeline may prefer a denormalised document/column model, while account balances require a strongly controlled transactional design. MongoDB recommends modelling around application access patterns and choosing embedding versus references accordingly (MongoDB data modelling).
2. Differentiate between Relational (SQL) databases and NoSQL databases based on data model, scalability, consistency, schema, and applications.
| Dimension | Relational/SQL | NoSQL |
|---|---|---|
| Data model | Tables, rows, columns, primary/foreign keys, joins | Key-value, document, wide-column, or graph; often denormalised |
| Schema | Schema-first; constraints and types are central | Flexible/schema-on-read or per-document schema; validation remains possible |
| Scaling | Strong vertical scaling; distributed sharding/replication is possible but more involved | Horizontal partitioning and replication are usually fundamental design features |
| Consistency/transactions | ACID transactions and strong consistency are common, ideal for integrity-critical updates | Varies: strong per record, tunable, eventual, or transactional within a bounded scope; distributed transactions may be limited/costly |
| Joins | Rich declarative joins and ad-hoc relational queries | Often avoid joins; duplicate/materialise data for known query paths |
| Access pattern | Flexible queries over normalised data | Model is selected around high-value predictable queries and partition keys |
| Availability | May trade availability during a partition to preserve consistency, depending on design | Many distributed systems prioritise availability and partition tolerance for selected workloads, with explicit consistency trade-offs |
| Best applications | Ledger, payroll, inventory, order commit, complex reporting | Product profiles, sessions, event streams, social graphs, IoT, content, high-scale carts |
The CAP theorem is not a licence to ignore correctness: under a network partition, a distributed system cannot simultaneously guarantee all three of linearizable consistency, availability, and partition tolerance. Actual systems offer nuanced/tunable choices rather than a simple “SQL consistent, NoSQL inconsistent” rule.
Example: an order database may atomically reserve stock and record payment in SQL; a document store may hold a denormalised product page; a wide-column store may keep device readings keyed by device/time; a graph database may traverse social connections. Hybrid “polyglot persistence” is often safer than migrating every workload. NoSQL’s flexibility and scale reduce some costs but shift responsibility to application code for data duplication, invariant enforcement, migrations, and operational observability.
3. Explain the four major NoSQL database types—Key-Value Store, Document Store, Column Family Store, and Graph Database—with suitable examples and applications.
| Type | Key-value logic/data model | Example | Suitable applications |
|---|---|---|---|
| Key-value | key -> opaque value; lookup, put, update, delete by key. The database need not understand fields inside the value. | Redis, Amazon DynamoDB | Sessions, caches, shopping carts, feature flags, counters. |
| Document | document_id -> JSON/BSON document; fields/nested arrays can be indexed and queried. Related bounded data can be embedded. | MongoDB, Couchbase | Product catalogues, profiles, content, reviews, order aggregates. |
| Column-family/wide-column | Rows are addressed by a row/partition key; columns are grouped into families and may be sparse. Efficient access requires a designed partition/clustering key. | Apache Cassandra, HBase | Time series, IoT, event history, high-write telemetry, large sparse tables. |
| Graph | Vertices/nodes and edges/relationships, both with properties; traversals follow relationships rather than joining many tables. | Neo4j, JanusGraph | Fraud rings, social networks, recommendations, dependency/network analysis. |
Key-value example: cart:U17 -> {P4:2,P9:1,expires:...}. A request with key cart:U17 returns the whole value; a query such as “all carts containing P4” is not naturally supported without an index or secondary structure.
Document example:
{"_id":"P4", "name":"phone", "price":399,
"reviews":[{"user":"U1","stars":5}], "variants":[...]}
Embedding bounded reviews gives one aggregate read, while unbounded/highly independent reviews should be referenced or stored separately. MongoDB documents have a 16 MiB limit, so large media should be referenced rather than embedded (MongoDB embedding guidance).
Wide-column example: key (device_id, day) with clustering column event_time and columns temperature, battery; a query supplies the partition key and scans a time range. Cassandra maps partition keys to token ranges and replicates partitions across nodes (Cassandra architecture). Graph example: (User A)-[BOUGHT]->(Product X) and (User A)-[FOLLOWS]->(User B); a recommendation traverses related users/products. Each type scales differently; the key/partition design determines performance.
4. Describe various NoSQL Data Architecture Patterns. Explain how these patterns address the challenges of handling Big Data.
- Key-value pattern: Choose a stable, high-cardinality key and distribute
key → valueacross partitions. It gives predictable O(1)-style point access and horizontal scale for sessions, carts, and counters. It sacrifices arbitrary field queries unless secondary indexes/materialised views are added. - Document aggregate pattern: Store an application aggregate—such as product plus bounded variants—in one document. One read and a single-document update are efficient; schema evolution is flexible. Unbounded child data is referenced or separated to avoid oversized/hot documents.
- Wide-column/time-series pattern: Partition by entity and time bucket, cluster by timestamp, and write append-heavy sparse columns. Bucketing prevents an ever-growing/hot partition; replication and sequential reads support telemetry. The query must provide the partition key; cross-partition scans are expensive.
- Graph pattern: Represent entities as vertices and relationships as edges. It makes multi-hop queries (friends-of-friends, fraud paths) natural, but graph traversals and global analytics need careful partitioning and may not scale like independent key lookups.
- Sharding/partitioning: Hash keys distribute load evenly; range or geo/time partitioning supports locality and range scans. A good key avoids hotspots and includes the access pattern.
- Replication and consistency pattern: Replicate data synchronously or asynchronously; choose quorum/tunable consistency and conflict resolution according to business risk. Cassandra documents partition-based replication and tunable consistency (Cassandra guarantees).
- Denormalisation/materialised-view pattern: Duplicate data in read-optimised shapes, update it through events or controlled writes, and accept reconciliation complexity. This replaces distributed joins with predictable reads.
- Polyglot persistence: Use a document store for content, wide-column store for telemetry, graph store for relationships, and SQL for financial invariants, connected by an event/data pipeline.
- CQRS/event-sourcing variation: separate write/event history from read projections; rebuild projections when logic changes, but manage ordering, idempotence, and eventual consistency.
These patterns address volume and velocity through partitioning, variety through flexible models, and availability through replication. They do not remove design work: the architect must document keys, consistency, failure recovery, retention, indexes, schema versions, and how duplicate projections are repaired.
5. Explain the concept of Shared-Nothing Architecture. How does it improve scalability and fault tolerance in NoSQL systems?
In a shared-nothing architecture, each node owns independent CPU, memory, and storage; nodes do not depend on a common disk or memory bus. Data is partitioned across nodes, usually by hashing or range-partitioning a key, and replicas are placed on other nodes.
client/router
|
+--------------+--------------+
| | |
Node A owns P0 Node B owns P1 Node C owns P2
disk/memory disk/memory disk/memory
| \\ replicas | \\ replicas | \\ replicas
+-------network coordination-------+
Scalability. If a key-space is partitioned among N comparable nodes, storage and aggregate throughput can grow approximately with N until coordination, skew, or network limits dominate. Adding a node rebalances only some partitions rather than replacing a central machine. Local CPU, memory, and disk prevent a shared SAN or memory bus from becoming the bottleneck.
Fault tolerance. A node failure loses its local copy but not necessarily the partition: replicas on other nodes can serve reads or accept writes according to the consistency policy. Membership/ring metadata detects failure, and re-replication restores the desired factor. For example, Cassandra uses partition keys and replication across nodes/data centres; the system is designed so a coordinator routes a request to replicas (Cassandra architecture).
Costs. Cross-partition joins and transactions require network coordination, so they are slower and more failure-prone. A poor partition key creates a hotspot, while a very large partition causes storage and repair problems. Rebalancing generates network/disk load; replicas consume capacity; eventual consistency may expose temporary stale values. Therefore design queries to be single-partition where possible, choose replication and quorum levels deliberately, use idempotent retries, and monitor skew/repair. Shared-nothing improves scale and availability by removing central hardware dependencies, not by making distributed coordination free.
6. Compare the Master–Slave and Peer-to-Peer distribution models used in NoSQL systems. Discuss their advantages, disadvantages, and suitable application areas.
| Aspect | Master–slave (primary–replica) | Peer-to-peer (leaderless/equal nodes) |
|---|---|---|
| Write path | One master/primary accepts writes; replicas copy them synchronously or asynchronously | Coordinator routes to a replica set; multiple peers can accept writes |
| Read path | Primary or replicas, depending on consistency/read preference | One or more replicas; quorum/tunable consistency may be used |
| Failover | Elect/promote a replica when master fails | No single master; membership/ring/routing continues around failed peer |
| Strength | Simple ordering, central authority, strong primary semantics, straightforward transactions | No single bottleneck, high write availability, natural horizontal and multi-site scale |
| Weakness | Primary can bottleneck or be a single failure/failover point; replication lag and stale reads | Conflict resolution, quorum latency, repair, duplicate writes, and consistency reasoning are harder |
| Suitable use | Document primary/replicas, read-heavy catalogue, workloads needing ordered writes | Globally distributed events, telemetry, high-write availability, partition-tolerant services |
A master–slave system can scale reads with replicas, but writes converge at the master. If the master fails, election/failover causes a period of unavailability or may risk conflicting writes unless fencing and logs are correct. Synchronous replicas improve durability but increase write latency; asynchronous replicas reduce latency but can lose recent writes on catastrophic failure.
In peer-to-peer systems, any node may coordinate a request and replicas use quorum rules. If replication factor is R=3, a write quorum W=2 and read quorum Q=2 overlap (W+Q>R) under the relevant assumptions, improving read-your-latest behaviour; availability and latency fall if too many replicas are unreachable. Eventual reconciliation, hinted handoff/read repair/repair jobs, timestamps or conflict-free rules may be needed. Cassandra is a peer-oriented distributed system using partitioning and replication; its guarantees explain that consistency and availability choices depend on the operation (Cassandra guarantees).
Neither model is universally superior. Choose master–slave for simpler invariants and read scaling; choose peer-to-peer for multi-region availability and sustained writes when the application can tolerate/tune consistency. Explicitly state failover, conflict, and stale-read semantics in an exam answer.
7. Analyze the suitability of different NoSQL databases for the following applications: (i) Social Media, (ii) Banking, (iii) E-commerce, and (iv) IoT Systems. Justify your answer.
| Application | Suitable primary model | Design and justification | Caution |
|---|---|---|---|
| Social media | Graph store plus document/column store | Graph nodes/edges represent follows, likes, and communities; traversals support mutual friends and recommendations. Documents hold posts/profiles; wide-column timelines can be precomputed by user/time. | Global graph traversals, hot celebrity keys, privacy, and feed fan-out need sharding/caching/materialised feeds. |
| Banking | Relational/strongly consistent store as system of record; NoSQL selectively | Balances, transfers, and ledger invariants require ACID controls, serialisation, audit, and regulatory correctness. Key-value/document stores can serve immutable statements, session/risk features, or high-volume transaction events after the authoritative commit. | Do not use eventual-consistent cache as the ledger; define idempotency, reconciliation, encryption, and exactly-once business effects. |
| E-commerce | Document store for catalogue/orders; key-value for carts/sessions; optional graph for recommendations | Product variants and attributes vary, so documents fit catalogue reads. cart:user_id gives a direct key-value lookup. Orders may be immutable documents with state transitions; recommendations may use graph/feature stores. | Inventory reservation/payment needs transactional coordination; denormalised copies and search indexes need update/reconciliation. |
| IoT | Wide-column/time-series store (Cassandra/HBase) plus object/lake archive | Partition by device_id and time bucket, cluster by event time, and replicate across nodes. It sustains append-heavy writes and range reads for a device. | Avoid unbounded partitions and hotspot keys; retain/downsample old data; use streaming for alert latency. |
Decision process. First identify the dominant access path and correctness requirement, then choose partition key, replication, consistency, retention, and query indexes. A social post lookup is a key lookup, but “friends who bought products liked by friends” is a graph traversal. A bank may use Cassandra for immutable telemetry but SQL for money movement. A retailer can use MongoDB for flexible product documents and Redis-like key-value access for carts; MongoDB advises embedding bounded related data and referencing independently queried/high-cardinality data (modelling guidance).
Thus “NoSQL database” is not one product. Polyglot persistence is often the justified answer, with an event pipeline and clear ownership of each fact. Benchmark representative reads/writes and failure scenarios before selection.
8. Evaluate the advantages and limitations of NoSQL databases over traditional relational databases for handling Big Data applications.
Advantages
- Horizontal partitioning distributes storage and throughput across nodes and supports elastic growth.
- Flexible document/column schemas handle evolving, nested, sparse, and semi-structured records.
- Denormalised aggregates and key-based access provide predictable high throughput without distributed joins.
- Replication and multi-region designs can improve availability and locality.
- Models match specialised workloads: graph traversal, time series, carts, sessions, and event ingestion.
- Many systems support commodity/cloud deployment and developer-friendly JSON or simple APIs.
Limitations
- Denormalisation duplicates facts; updates need application logic, events, transactions, or repair jobs.
- Arbitrary ad-hoc joins and global aggregates are harder; a query that omits the partition key may require expensive scatter-gather.
- Consistency varies. Eventual reads can expose stale inventory/profile data; quorum and distributed transactions add latency and complexity.
- Schema flexibility can become poor quality without validation, versioning, catalogue, and migration discipline.
- Indexes, large partitions, hot keys, compaction, repair, and rebalancing require specialised operations.
- NoSQL does not automatically provide ACID multi-record invariants, SQL maturity, reporting tools, or regulatory audit semantics.
- Vendor/product differences make portability and skills a concern; replication and storage overhead still exist.
Evaluation example. A document catalogue can scale product reads and tolerate new attributes better than a normalised schema. However, stock decrement must not be based on a stale document; use a transactional authority/conditional update and reconcile projections. Cassandra’s documentation notes that partition keys determine placement and efficient queries supply them; it avoids cross-partition joins by design (Cassandra overview).
The correct conclusion is workload-dependent: NoSQL is preferable for distributed high-volume/high-availability access patterns, while relational databases remain preferable for strongly consistent transactions, complex joins, and integrity constraints. Most enterprises combine them, selecting by latency, consistency, query shape, data evolution, cost, and compliance.
9. Case Study: A multinational online shopping company manages millions of products, customer reviews, shopping carts, and recommendation data. Recommend a suitable NoSQL database architecture and justify your choice.
Requirements. Catalogue reads must be fast and globally available; attributes vary by category; reviews are append-heavy and independently paged; carts need low-latency per-user updates and expiry; recommendations need relationship/feature queries; traffic is bursty; regional failure must not lose committed cart/order intent; and personal data needs isolation and deletion controls.
Recommended polyglot architecture
catalog/admin ---> document store (product, variants, bounded attributes)
reviews ---------> wide-column/document review store (product_id,time,review_id)
carts -----------> key-value store: cart:{user_id} -> cart aggregate + TTL
clicks/orders ----> event bus -> lake/batch features -> recommendation store
users/products ---> graph or offline feature graph -> top-N recommendations
|
API/query services + cache + search index
A document store is the primary catalogue model: one product document contains bounded variants and category-specific attributes. Embed data read/updated together; reference reviews because they are high-cardinality, independently paged, and unbounded. MongoDB’s guidance explicitly recommends modelling around access patterns and using references when children are queried or changed independently (MongoDB best practices). Store images in object storage and keep URLs/metadata in the document.
Use a key-value store for cart:{customer_id} with conditional version/compare-and-set updates, TTL for abandoned carts, and regional replication. Checkout must revalidate price and inventory in a transactional authority; a cart is not proof of stock. Use a wide-column store for review events/ratings keyed by (product_id, time_bucket) and a moderation status. Build recommendations offline from views, purchases, and product similarity; serve top-N lists by (customer_id, region) from a key-value/document feature store. A graph database is justified if online multi-hop relationship queries are required; otherwise offline graph computation plus key-value serving is simpler.
Data flow and trade-offs. APIs write authoritative events; an outbox/event stream updates denormalised projections idempotently. Read services use the model matching each query. Shard by product/customer ID, avoid hot keys, replicate across regions, encrypt PII, and provide per-tenant access/audit. Eventual consistency is acceptable for recommendations and review counts, but not silently for payment/inventory. The architecture gains scale and query performance at the cost of duplicated data, asynchronous lag, multiple operational technologies, reconciliation, and more complex disaster recovery. A single document database is acceptable for a smaller deployment; the multinational workload justifies polyglot NoSQL only with measured access patterns and operational capacity.
10. Application-Based Question: Design a NoSQL solution for a smart healthcare system that stores patient records, medical images, sensor data, and doctor’s notes. Select an appropriate NoSQL database type and explain your design with justification.
Choice. Use a document store as the patient-record and notes system, complemented by a wide-column/time-series store for high-rate sensors and encrypted object storage for medical images. If the exam requires one primary NoSQL type, choose a document store because patient records and doctor notes are nested, evolve over time, and are retrieved as patient-centred aggregates. Do not put large image binaries directly into ordinary documents; store them in encrypted object/blob storage and reference them.
EHR/doctor app ---> document DB
patient_id -> demographics,
encounters, medications, notes,
consent, version/audit metadata
sensor gateway ---> wide-column DB: (patient/device, day, event_time)
images ----------> encrypted object store: image_id/path/checksum
|
metadata links in document DB
|
authorised clinical API / analytics lake / alerts
Example document (conceptual):
{
"patient_id":"P17", "version":42,
"consents":[{"purpose":"care","status":"active"}],
"encounters":[{"id":"E8","time":"...","diagnoses":["..."],
"notes":[{"author":"D4","text_ref":"N91","signed":true}]}],
"images":[{"image_id":"I3","object_uri":"...","sha256":"..."}]
}
Embed bounded data that is read together, but store unbounded encounters/notes as separate documents keyed by (patient_id, encounter_id) if the patient aggregate would grow without bound. Maintain immutable signed note versions; never overwrite the legal record invisibly. The document store’s flexible validation allows new clinical fields while schema versions and terminology mappings preserve quality.
For sensors, partition by device_id/patient_id and time bucket, cluster by timestamp, and replicate across nodes. This supports append-heavy writes and range queries such as “last 24 hours,” while downsampling/archive controls partition and storage growth. Cassandra describes partition-key placement and replication across a distributed cluster (architecture overview). Images are encrypted at rest and in transit; access uses short-lived authorised URLs, not public paths.
Security and trade-offs. Use patient pseudonyms, RBAC/ABAC, consent checks, encryption, audit trails, retention/deletion policies, backups, regional residency, and a break-glass procedure. Strong conditional updates protect note signing and care-plan versions; sensor dashboards may tolerate eventual consistency, but clinical alerts need a validated stream processor. A NoSQL system improves scale and schema agility but does not itself guarantee clinical correctness, interoperability, privacy, or ACID cross-store transactions. Use an outbox/event log and reconciliation; keep the authoritative clinical/legal ownership explicit.
Compact authoritative sources
- Apache Hadoop modules and HDFS Architecture
- Apache Hadoop MapReduce Tutorial
- HDFS High Availability
- Apache Hive, Pig, HBase, Sqoop, and Flume
- Leskovec, Rajaraman & Ullman, Mining of Massive Datasets, Chapter 2
- Apache Cassandra architecture and guarantees
- MongoDB data modelling and best practices