|
|
关于MongoDB,我们能看到的资料,基本都是在指导大家如何使用MongoDB,但是,MongoDB内部是如何运作的,资料不是很多。) ~' w Q7 O" D, v0 T
% M3 A6 {$ }2 Y! r" n9 e: S
阅读使用手册,会有很多疑惑之处。例如,有人说,MongoDB 等同于分布式的 MySQL。它把一个Table ,按 row,分割成多个Shards,分别存放在不同的 Servers 上。这种说法是否正确?) O- \ S; a) I: V
8 [: r Z( D% Z
不深入了解 MongoDB 的内部结构,就无法透彻地回答类似问题。这个系列文章,就来和大家探讨MongoDB的内部的工作方式。
5 h/ G, h2 e1 C
$ f' b, U- ]9 I' _% F
; x# Y5 u" `! a& Z
?) ^+ q Q9 e+ [" w! o& g3 B5 ]图1-1 MongoDB架构图
1 J- |! Z! v- [4 l7 ~) C& W& \: P* X0 h2 f* q9 _
MongoDB 通常运行在一个服务器集群上,而不是一个单机。图1-1,描述了一个MongoDB集群的基本组成部分,包括若干shards,至少一个config server,至少一个routing servers(又称 mongos)。
9 ]. j$ a6 |! K! \' I5 m- [- s q
Shards) J; L9 R) M; {* z, i! F
) F4 J6 E; f* I* C% U+ Z
MongoDB的最基本的数据单元,叫document,类似于关系式数据库中的行 row。一系列documents,组成了一个collection,相当于关系式数据库中的table。当一个 collection 数据量太大时,可以把该collection按documents切分,分成多个数据块,每个数据块叫做一个chunk,多个chunks聚集在一起,组成了一个shard。% h( X" |- w. I
0 Z2 v0 e8 \, f6 s/ y/ C/ G Sharding 的意义,不仅保障了数据库的扩容(scalability),同时也保障了系统的负载均衡(load balance)。
2 Y1 @# d- x {: H; ~
* }; i' g2 u, A; `% n 每一个shard存储在一个物理服务器(server)上。Server上运行着mongod进程,通过这个进程,对shard中的数据进行操作,主要是增删改查。7 B0 ~) r' m2 G& A, ~ U
6 Y1 A1 R' e9 _. h2 T9 d- V, f 如果系统中的每个shard,只存储了一份数据,没有备份,那么当这个shard所在的server挂了,数据就丢失了。在生产环境中,为了保证数据不丢失,为了提高系统的可用性(availability),每一个shard被存储多份,每个备份所在的servers,组成了一个replica set。
1 T0 M% w' k. H; {
+ o. E8 k* l% E* v2 {2 R. A/ dShard keys; g- \$ g) s% `8 H/ |8 ^ f$ l
( w5 {; a) }+ v! H# T$ ~" L9 X1 L 为了把collection切分成不同的chunks,从而存放到不同的shards中,我们需要制定一个切分的方式。& ?1 W5 f- K4 E& G( O
* t/ ?! ?+ F" {" x: V- o. o2 C 如前所述,在 MongoDB 数据库中,一个表collection由多个行 documents 组成,而每个 document,有多个属性 fields。同一个 collection 中的不同的 documents,可能会有不同的 fields。例如,有个 collection 叫 Media,包含两条 documents,) q. l8 S5 D( \0 s$ l7 A1 |
8 {# {+ @3 k; s- C8 s: W1 k" g. C
{6 W& x/ W3 m4 K. a9 f9 x3 X& D! b
"ISBN": "987-30-3652-5130-82",( B( `0 ]# b( O% M/ L6 t! j
"Type": "CD",
% L: I2 N$ x6 P "Author": "Nirvana",
* p+ ]6 r. Y2 G# C3 ~* n0 y "Title": "Nevermind",( c/ C' d' o0 J
"Genre": "Grunge",
7 F% o1 W @* z9 [& r- v7 k "Releasedate": "1991.09.24"," z8 B) C( n8 l3 {0 {8 d" T
"Tracklist": [
H! B5 t2 A. M( @4 i {) D) ^( P* @5 y# W7 D& U
"Track" : "1",
! k1 ]2 X) q7 }( V "Title" : "Smells like teen spirit",9 s2 d* o; s1 r# _; N
"Length" : "5:02"
6 h$ J S. k- b8 m: u) Y+ S" t },
2 }( F p+ `) N {" t, }) t/ v+ s* U6 d. `, L
"Track" : "2",
/ ]1 J- D" d/ r7 j "Title" : "In Bloom",
" Q+ b4 ]' a$ [7 J3 F, n! q "Length" : "4:15"
7 \* [0 v" |! }, t3 k0 }4 x }
# a2 j9 v" k! {9 o$ S) q5 S/ [ ]
7 L" h$ U9 y* C1 c0 {}
" Q( E' H* C8 X4 N
, U+ x- I' V1 {6 g/ Y% d- p+ Q7 [{
% B" R6 c4 ~$ e; j4 B- Q3 }( X "ISBN": "987-1-4302-3051-9",; G# ?" P3 {" j9 J* `0 x
"Type": "Book",
$ e! S" _# ^' ^- e+ T "Title": "Definite Guide to MongoDB: The NoSQL Database",
, V# s5 N. |* b1 l "Publisher": "Apress",3 _3 o, ]; U! {" a
"Author": " Eelco Plugge",
/ U7 `, o; ]# e# d, ] "Releasedate": "2011.06.09"
C, @/ h8 f# {}
; | O, o- q, r
' S, `/ p9 A h# h3 n 假如,在同一个 collection 中的所有 document,都包含某个共同的 field,例如前例中的“ISBN”,那么我们就可以按照这个 field 的值,来分割 collection。这个 field 的值,又称为 shard key。
) N/ J) A5 i, Y0 l7 v- M: r
1 }$ W: w+ g; I% P6 j 在选择shard key的时候,一定要确保这个key能够把collection均匀地切分成很多chunks。
. f* J* a/ [1 L) f0 m8 `- m( X" q6 q) Z& F! j: W P
例如,如果我们选择“author”作为shard key,如果有大量的作者是重名的,那么就会有大量的数据聚集在同一个chunk中。当然,假设很少有作者同名同姓,那么“author”也可以作为一个shard key。换句话说,shard key 的选择,与使用场景密切相关。
! `8 U5 q! i+ _
: Z" ^( s; R( G9 Y. b 很多情况下,无论选择哪一个单一的 field 作为shard key,都无法均匀分割 collection。在这种情况下,我们可以考虑,用多个 fields,构成一个复合的shard key。
) B2 Y5 V7 O" O* j7 K1 X. t" }) ?* l/ E
; d( R1 T8 M" l- x6 h& I/ G 延续前例,假如有很多作者同名同姓,他们都叫“王二”。用 author 作为 shard key,显然无法均匀切割 collection。这时我们可以加上release-date,组成name-date的复合 shard key,例如“王二 2011”。
; R* h8 t# j2 I) i3 u; G. V) |: I# N+ O) ? @ ~0 ~
Chunks1 G/ @4 n$ m% s6 W4 `
0 ^$ ^7 }$ K1 R4 I# H0 n
MongoDB按 shard key,把 collection切割成若干 chunks。每个 chunk 的数据结构,是一个三元组,{collection,minKey,maxKey},如图1-2 所示。$ D; S( F: K. i
+ O/ s/ p. o, K3 _7 m9 D
( k; {: Z6 K2 O9 [图1-2 chunk的三元组
8 X0 S2 F7 ?7 @/ E3 c& O2 H7 i9 c+ T( j0 d8 P" ]2 G6 L. t% ]9 H! k
其中,collection 是数据库中某一个表的名称,而 minKey 和 maxKey 是 shard key的范围。每一个 document 的shard key 的值,决定了这条document应该存放在哪个chunk中。6 a& i/ F! c0 k7 [
9 J* C7 A M0 z0 ?, S6 I
如果两条 documents 的 shard keys 的值很接近,这两条 documents 很可能被存放在同一个 chunk 中。
- Q2 G/ f; F8 Q7 @& R: p* a
* d" ^2 p: ^+ F( b! u, w) d- a Shard key 的值的顺序,决定了 document 存放的 chunk。在 MongoDB 的文献中,这种切割 collection 的方式,称为order-preserving。+ Y1 P* f- [, X
' ]3 Y) k4 g' F/ p, x. q 一个 chunk最多能够存储64MB的数据。 当某个chunk存储的 documents包含的数据量,接近这个阈值时,一个chunk会被切分成两个新的chunks。& |: U3 [2 c* i U* s
) z; r' N7 p7 i6 \ } 当一个shard存储了过多的chunks,这个shard中的某些chunks会被迁移到其它 shard中。
) C1 k1 v, P0 _9 { u" [# N0 p) ~# ] [$ }, w1 l8 g
这里有个问题,假如某一条 document 包含的数据量很大,超过 64MB,一个 chunk 存放不下,怎么办?在后续章节介绍 GridFS 时,我们会详细讨论。2 e( }* U, ~) t: E7 {& a
0 G# ]4 A q8 [Replica set
- E' y( R: N% l7 b
7 q3 H4 s" P# M9 y 在生产环境中,为了保证数据不丢失,为了提高系统的可用性(availability),每一个shard被存储多份,每个备份所在的servers,组成了一个replica set。- d, z* G# A6 o9 ~# U: L
9 F* h* A7 x' x2 ?
这个replica set包括一个primary DB和多个secondary DBs。为了数据的一致性,所有的修改(insert / update / deletes) 请求都交给primary处理。处理结束之后,再异步地备份到其他secondary中。
) b7 w* {3 Y+ n: g5 Z5 p/ {
! o2 X- s+ T- `7 }' K" y% y8 i Primary DB由replica set中的所有servers,共同选举产生。当这个primaryDB server出错的时候,可以从replica set中重新选举一个新的primaryDB,从而避免了单点故障。" V8 Z* A4 U) ~
# O, K! @$ D4 ~. c9 P Replica set的选举策略和数据同步机制,确保了系统的数据的一致性。后文详述。' r7 s3 @' t9 s; \, B# N
4 r# \. u8 J6 i
Config Server2 H+ k, L$ m% v1 G. M
# a1 f, `* A4 e8 E. S2 Y Config servers用于存储MongoDB集群的元数据 metadata,这些元数据包括如下两个部分,每一个shard server包括哪些chunks,每个chunk存储了哪些 collections 的哪些 documents。5 k4 w9 }& w* P% o7 L8 d, Q, s9 ~
5 B+ B" q U; G8 C( n) `4 U7 ^
每一个config server都包括了MongoDB中所有chunk的信息。
- e; f( n g0 z$ ~0 H1 m4 F( z! d+ X$ {/ ^# z, }
Config server也需要 replication。但是有趣的是,config server 采用了自己独特的replication模式,而没有沿用 replica set。, X" I3 w& i+ X
. F) |8 k" X( Y+ d5 B
如果任何一台config server挂了,整个 config server 集群中,其它 config server变成只读状态。这样做的原因,是避免在系统不稳定的情况下,冒然对元数据做任何改动,导致在不同的 config servers 中,出现元数据不一致的情况。$ \9 w# F! K) @" X% ?, c
+ M: f/ M) p" y, Y# { F MongoDB的官方文档建议,配置3个config servers比较合适,既提供了足够的安全性,又避免了更多的config servers实例之间的数据同步,引起的元数据不一致的麻烦。
6 B; q; ?0 n: N/ m& {; h
1 Y' P8 H8 z* p/ OMongos9 x( W( n6 U" ]2 X3 [
. @7 f& v% x2 b) O
用户使用MongoDB 时,用户的操作请求,全部由mongos来转发。
+ W! V$ T& b/ b2 w0 D$ q' p0 x4 A4 T$ E) @
当 mongos 接收到用户请求时,它先查询 config server,找到存放相应数据的shard servers。然后把用户请求,转发到这些 shard servers。当这些 shard servers完成操作后,它们把结果分别返回给 mongos。而当 mongos 汇总了所有的结果后,它把结果返回给用户。+ ~1 R4 G$ o2 k: n; M
, b$ i! q; I+ L( X
Mongos每次启动的时候,都要到config servers中读取元数据,并缓存在本地。每当 config server中的元数据有改动,它都会通知所有的mongos。& \7 _: P) o5 D$ f3 h2 }
' ?5 R ?& g8 W5 ~ H Mongos之间,不存在彼此协同工作的问题。因此,MongoDB所需要配置的mongos server的数量,没有限制。: }3 o* T1 v$ \1 q! L- K8 `
+ [( b% \8 W8 ?0 p3 V 通过以上的介绍,我们对每个组成部分都有了基本的了解,但是涉及到工作的细节,我们尚有诸多疑问,例如,一个chunk的数据太大,如何切分?一个shard数据太多,如何迁移?在replica set中,如何选择primary?server挂了,怎么进行故障恢复?接下来的章节,我们逐个回答这些问题。5 A5 o4 n9 f6 j! u0 i
[1 a. V" z: f0 m7 C+ T
7 V+ Y, v5 F# d7 ]$ {
Reference,) p6 a$ O) w8 @% d
9 m% e V- b' U
[0] Architectural Overview
s4 c/ L- y" T$ e) vhttp://www.mongodb.org/display/DOCS/Sharding+Introduction
# E/ M8 X K$ C, b0 G( l( p0 w |
评分
-
查看全部评分
|