Termín Big Data je v roce 2012 vsude — na konferencích, v médiích, v nabidkach vendorů. Ale co to prakticky znamená pro českou firmu, která zpracovává desítky az stovky gigabytu dat denně? Apache Hadoop muze byt odpověď.
Co je Big Data a kdy ho skutecne potřebujete¶
Gartner definuje Big Data pomoci tri V: Volume (objem), Velocity (rychlost) a Variety (rozmanitost). Pokud vaše data splňují alespoň dvě z těchto kritérií a stávající relační databaze nestačí, je čas se podívat na alternativy.
Typicke scénáře, kde se Hadoop vyplatí:
- Analýza logu — webove servery, aplikacni servery, síťové prvky generují gigabajty logu denně. SQL dotazy nad takovou tabulkou trvají hodiny.
- ETL pro datový sklad — transformace a čištění dat pred nahranim do Oracle nebo SQL Server data warehouse
- Analýza zákaznického chování — clickstream data z e-shopu, telekomunikační CDR záznamy
- Archivace a vyhledávání — staré dokumenty, emaily, scany — nestrukturovaná data, kde fulltext nestačí
Pokud zpracováváte méně nez 10 GB dat a vaše dotazy běží v rozumném case, Hadoop pravdepodobne nepotřebujete. Relační databaze s dobrými indexy a materializovanymi pohledy jsou pro mensi objemy efektivnější.
Architektura Hadoop clusteru¶
Hadoop se skládá ze dvou klíčových komponent:
HDFS (Hadoop Distributed File System) — distribuovaný souborový system, který data replikuje na více nodu (výchozí replikační faktor je 3). Data jsou rozdelena na bloky o velikosti 64 MB (v Hadoop 2.x typicky 128 MB) a distribuovaná napříč clusterem.
MapReduce — výpočetní framework, který zpracovává data paralelně na všech nodech, kde data leží. Misto přesunu dat k vypoctu se výpočet přesouvá k datum — to je klíčový princip.
Minimalni produkci cluster pro enterprise použití:
- NameNode — 1 server, 64 GB RAM, RAID 1, ridici uzel HDFS
- Secondary NameNode — 1 server, backup metadat
- DataNode / TaskTracker — minimálně 4 servery, kazdy 32 GB RAM, 4-12 disku bez RAID (HDFS replikuje sam), 8+ jader CPU
- Edge node — 1 server pro klientsky pristup, Hive, Pig, import/export dat
Distribuce: Cloudera vs Apache vs Hortonworks¶
Ciste Apache Hadoop je možné provozovat, ale pro enterprise nasazeni doporučujeme jednu z komerčních distribuci:
Cloudera CDH 4 — nejrozšířenější enterprise distribuce. Obsahuje Hadoop, Hive, HBase, Pig, Oozie a Cloudera Manager pro spravu clusteru. Komerční podpora a certifikace. Pro vetsinu našich klientu doporučujeme tuto variantu.
Hortonworks HDP 1.x — plne open-source distribuce. Bez vlastního management nastroje (používá Apache Ambari). Vhodná pro firmy s vlastním Hadoop know-how.
MapR — nahrazuje HDFS vlastním vysokovykonnym souborovym systémem. Zajímavá volba pro nízkou latenci, ale vendor lock-in.
Hive: SQL nad Hadoop¶
Pro analytické dotazy nad daty v HDFS je Apache Hive ideální volbou. Hive umožňuje pisat dotazy v jazyce podobném SQL (HiveQL), které se internene převedou na MapReduce joby.
– Příklad: analýza pristupu k webu za poslední mesic
CREATE EXTERNAL TABLE access_log (
ip STRING,
request_time STRING,
method STRING,
url STRING,
status INT,
bytes BIGINT
) ROW FORMAT DELIMITED FIELDS TERMINATED BY ‘\t’
LOCATION ‘/data/logs/access/’;
SELECT url, COUNT(*) as hits, SUM(bytes) as total_bytes
FROM access_log
WHERE status = 200
GROUP BY url
ORDER BY hits DESC
LIMIT 100;
Tento dotaz zpracuje terabajty logu paralelně na celém clusteru. Na relační databazi by trvat hodiny — na Hadoop clusteru s 8 nody minuty.
Integrace s existujícím ekosystémem¶
Hadoop není náhrada za Oracle nebo SQL Server. Je to doplněk. Typicky workflow:
- Sqoop — import dat z relační databaze do HDFS
- MapReduce / Hive — transformace a agregace
- Sqoop export — vysledky zpět do relační databaze pro BI nastroje
Pro real-time sber dat (logy, udalosti) pouzivame Apache Flume, který data streamuje primo do HDFS. Pro messaging mezi systemy je vhodný Apache Kafka (relativně nový projekt od LinkedIn, ale už stabilní).
Provozní aspekty¶
Hadoop cluster vyzaduje jiný pristup k provozu nez klasický aplikacni server:
- Monitoring — Ganglia nebo Nagios s Hadoop pluginy. Sledujte HDFS kapacitu, počet živý DataNodes, MapReduce queue a failed joby.
- Zálohování — HDFS má built-in replikaci, ale NameNode metadata jsou single point of failure. Zálohujte fsimage a edits log.
- Kapacitní plánování — data v HDFS rostou. Plan na 30-50 procent volneho místa. Přidávání DataNodes je snadné — Hadoop automaticky rebalancuje data.
- Bezpecnost — Hadoop ve výchozím nastaveni nemá autentizaci. Pro enterprise nasazeni zapnete Kerberos integraci.
Náklady a ROI¶
Hadoop běží na commodity hardware — to je jeho hlavní ekonomická vyhoda. Cluster s 8 DataNodes na standardních 2U serverech stojí řádově 1-2 miliony CZK včetně disku. Srovnatelný vykon na komerční MPP databazi (Teradata, Netezza) stojí nasobne více.
ROI se typicky projeví v těchto oblastech:
- Rychlejsi ETL procesy — z hodin na minuty
- Analýzy, které předtím nebyly možné (fulltext pres miliony dokumentu)
- Dlouhodobá archivace dat za zlomek ceny SAN úložiště
- Offload zateze z produkci databaze
Shrnuti¶
Hadoop není silver bullet, ale pro spravne use cases prinasi dramatické zlepšení. Zacnete s jedním konkrétním problémem — analýza logu nebo ETL offload — a rust clusteru podridte reálným potřebám. Cloudera CDH 4 je solidní volba pro ceske enterprise prostredi s dostupnou podporou. Klíčové je mit v tymu alespoň jednoho člověka, který Hadoop rozumí — at už interního nebo externího.
Potřebujete pomoc s implementací?
Naši experti vám pomohou s návrhem, implementací i provozem. Od architektury po produkci.
Kontaktujte nás