Kiến thức cốt lõi về hệ sinh thái dữ liệu lớn Big Data

1. Tổng quan về Big Data

Big Data được định nghĩa bởi các đặc trưng cơ bản giúp phân biệt với các hệ thống xử lý dữ liệu truyền thống:

  • Khối lượng (Volume): Sự bùng nổ dữ liệu từ các nguồn phi cấu trúc khiến quy mô lưu trữ vượt xa ngưỡng Terabyte, tiến dần đến Petabyte (PB) và Zettabyte (ZB).
  • Đa dạng (Variety): Dữ liệu không chỉ dừng lại ở dạng bảng (có cấu trúc) mà còn bao gồm nhật ký web, video, âm thanh, hình ảnh (phi cấu trúc) và dữ liệu bán cấu trúc.
  • Giá trị (Value): Mật độ giá trị trong dữ liệu lớn thường thấp do chứa nhiều nhiễu. Thách thức đặt ra là sử dụng thuật toán máy học và AI để trích xuất thông tin hữu ích từ biển dữ liệu thô.
  • Tốc độ (Velocity): Yêu cầu xử lý và phân tích dữ liệu phải diễn ra gần như thời gian thực để đảm bảo tính thời sự và hiệu quả kinh doanh.

Mối quan hệ giữa Big Data, Cloud Computing và IoT:

  • IoT: Nguồn phát sinh dữ liệu khổng lồ từ các cảm biến.
  • Cloud Computing: Cung cấp hạ tầng tài nguyên (tính toán, lưu trữ) linh hoạt để vận hành các hệ thống Big Data.
  • Big Data: Tập trung vào các kỹ thuật xử lý, phân tích và lưu trữ các luồng dữ liệu từ IoT trên nền tảng Cloud.

2. Hệ sinh thái Hadoop và cơ chế lưu trữ HDFS

Chế độ vận hành của Hadoop

  • Standalone Mode: Chế độ mặc định, chạy trên một node duy nhất, sử dụng hệ thống tệp cục bộ, chủ yếu dùng để gỡ lỗi.
  • Pseudo-Distributed Mode: Chạy trên một node nhưng mô phỏng các tiến trình phân tán (NameNode, DataNode) thông qua các tiến trình Java riêng biệt.
  • Fully Distributed Mode: Vận hành trên một cụm máy chủ thực tế, triển khai các dịch vụ trên nhiều node khác nhau.

Khái niệm Block trong HDFS

HDFS chia nhỏ các tệp tin thành các khối (Block) có kích thước cố định (mặc định 64MB hoặc 128MB) để quản lý.

  • Ưu điểm: Cho phép lưu trữ tệp lớn hơn dung lượng của một node đơn lẻ, đơn giản hóa quản lý Metadata và hỗ trợ sao lưu (Redundancy) hiệu quả.
  • Lưu ý: Nếu một tệp nhỏ hơn kích thước Block, nó sẽ không chiếm dụng toàn bộ không gian của Block đó trên đĩa cứng.

Cấu trúc NameNode và phục hồi dữ liệu

NameNode quản lý Metadata thông qua ba cấu trúc chính:

  • FsImage: Bản sao lưu trạng thái hệ thống tệp tại một thời điểm nhất định.
  • EditLog: Ghi lại mọi thao tác thay đổi (tạo, xóa, đổi tên) đối với tệp tin.
  • Metadata trong bộ nhớ: Chứa thông tin cấu trúc cây thư mục và ánh xạ các block.

Để đảm bảo an toàn, SecondaryNameNode sẽ định kỳ gộp EditLog vào FsImage để rút ngắn thời gian khởi động lại và dự phòng khi NameNode gặp sự cố.

3. Cơ sở dữ liệu HBase

HBase là cơ sở dữ liệu phi quan hệ (NoSQL) định hướng cột chạy trên nền tảng HDFS.

  • Mô hình dữ liệu: Dữ liệu được tổ chức theo Table, Row Key, Column Family, Column Qualifier và Timestamp.
  • Lưu trữ thưa thớt: Khác với DB truyền thống, HBase không lưu trữ các ô trống (null), giúp tiết kiệm không gian đáng kể.
  • Cơ chế ghi: Dữ liệu mới được ghi vào HLog (WAL)MemStore (bộ nhớ đệm). Khi MemStore đầy, nó sẽ được đẩy xuống đĩa dưới dạng StoreFile (HFile).

4. Hệ thống NoSQL và định lý CAP

Phân loại NoSQL

  • Key-Value: Redis, Riak.
  • Column-family: HBase, Cassandra.
  • Document: MongoDB, CouchDB.
  • Graph: Neo4j, OrientDB.

Định lý CAP

  • Consistency (C): Mọi node đều thấy dữ liệu giống nhau tại cùng một thời điểm.
  • Availability (A): Mọi yêu cầu đều nhận được phản hồi (thành công hoặc thất bại).
  • Partition Tolerance (P): Hệ thống vẫn hoạt động ngay cả khi mạng bị phân tách.

Trong thực tế, một hệ thống phân tán chỉ có thể đáp ứng tối ưu 2 trong 3 yếu tố này đồng thời.

5. Lập trình MapReduce

MapReduce xử lý dữ liệu qua hai giai đoạn chính:

  • Map: Chia nhỏ dữ liệu thành các cặp Key-Value.
  • Shuffle & Sort: Nhóm các cặp có cùng Key lại với nhau.
  • Reduce: Tổng hợp dữ liệu từ các nhóm Key để đưa ra kết quả cuối cùng.

Bài toán Join trong MapReduce

Giả sử cần Join hai bảng R(A, B) và S(B, C) qua thuộc tính B:

  1. Giai đoạn Map gán nhãn nguồn cho mỗi bản ghi: (b, (R, a))(b, (S, c)).
  2. Giai đoạn Reduce nhận danh sách các giá trị có cùng b, thực hiện tích Descartes giữa tập con từ R và tập con từ S để tạo ra (a, b, c).

6. Apache Spark và RDD

Spark cải tiến hiệu suất so với MapReduce nhờ cơ chế tính toán trong bộ nhớ (In-memory computing).

Khái niệm RDD (Resilient Distributed Dataset)

  • Immutable: Không thể thay đổi sau khi tạo.
  • Partitioned: Chia tách để xử lý song song.
  • Lazy Evaluation: Các phép biến đổi (Transformation) chỉ được thực thi khi có một hành động (Action) được gọi.
  • Lineage: Ghi lại vết lịch sử biến đổi để phục hồi dữ liệu khi có node bị lỗi.

Phân loại phụ thuộc (Dependency)

  • Narrow Dependency: Một Partition cha được dùng bởi tối đa một Partition con (ví dụ: Map, Filter). Không gây ra Shuffle.
  • Wide Dependency: Một Partition cha được dùng bởi nhiều Partition con (ví dụ: GroupByKey, Join). Gây ra Shuffle và chia tách các Stage trong DAG.

Mã nguồn ví dụ: WordCount trong Spark (Scala)


import org.apache.spark.{SparkConf, SparkContext}

object DistributedWordCount {
  def main(args: Array[String]): Unit = {
    val configuration = new SparkConf().setAppName("WordCountEngine").setMaster("local[*]")
    val context = new SparkContext(configuration)
    
    val rawData = context.textFile("hdfs:///data/input.txt")
    val counts = rawData.flatMap(content => content.split(" "))
                       .map(token => (token, 1))
                       .reduceByKey(_ + _)
    
    counts.saveAsTextFile("hdfs:///data/output")
    context.stop()
  }
}

7. Xử lý luồng với Apache Storm

Storm xử lý dữ liệu theo thời gian thực dưới dạng các Topology.

  • Spout: Nguồn dữ liệu đầu vào.
  • Bolt: Các đơn vị xử lý logic (lọc, tổng hợp, lưu trữ).

Chiến lược phân phối dữ liệu (Grouping)

  • Shuffle Grouping: Phân phối ngẫu nhiên và đồng đều các Tuple.
  • Fields Grouping: Các Tuple có cùng giá trị tại một trường chỉ định sẽ được gửi đến cùng một Task.
  • All Grouping: Sao chép và gửi Tuple đến tất cả các Task nhận.
  • Global Grouping: Toàn bộ luồng dữ liệu được gửi về một Task duy nhất có ID thấp nhất.

Mã nguồn ví dụ: WordCount Bolt (Java)


public class CounterBolt extends BaseBasicBolt {
    private Map<String, Long> _statistics = new HashMap<>();

    @Override
    public void execute(Tuple input, BasicOutputCollector collector) {
        String wordStr = input.getString(0);
        Long currentCount = _statistics.getOrDefault(wordStr, 0L);
        currentCount++;
        _statistics.put(wordStr, currentCount);
        collector.emit(new Values(wordStr, currentCount));
    }
}

Thẻ: hadoop Spark hbase Storm hdfs

Đăng vào ngày 5 tháng 8 lúc 04:38