1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20 package com.allanbank.mongodb.client.state;
21
22 import java.beans.PropertyChangeEvent;
23 import java.beans.PropertyChangeListener;
24 import java.beans.PropertyChangeSupport;
25 import java.net.InetSocketAddress;
26 import java.util.ArrayList;
27 import java.util.Arrays;
28 import java.util.Collections;
29 import java.util.List;
30 import java.util.concurrent.ConcurrentHashMap;
31 import java.util.concurrent.ConcurrentMap;
32 import java.util.concurrent.CopyOnWriteArrayList;
33
34 import com.allanbank.mongodb.MongoClientConfiguration;
35 import com.allanbank.mongodb.ReadPreference;
36 import com.allanbank.mongodb.Version;
37 import com.allanbank.mongodb.client.ClusterStats;
38 import com.allanbank.mongodb.client.ClusterType;
39 import com.allanbank.mongodb.client.Message;
40 import com.allanbank.mongodb.client.VersionRange;
41 import com.allanbank.mongodb.util.ServerNameUtils;
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63 public class Cluster implements ClusterStats {
64
65
66 public static final String SERVER_PROP = "server";
67
68
69 public static final String WRITABLE_PROP = "writable";
70
71
72 protected final MongoClientConfiguration myConfig;
73
74
75 protected final ConcurrentMap<String, Server> myServers;
76
77
78 protected VersionRange myServerVersionRange;
79
80
81 protected int mySmallestMaxBatchedWriteOperations;
82
83
84 protected long mySmallestMaxBsonObjectSize;
85
86
87 final PropertyChangeSupport myChangeSupport;
88
89
90 final ServerListener myListener;
91
92
93 final CopyOnWriteArrayList<Server> myNonWritableServers;
94
95
96 final CopyOnWriteArrayList<Server> myWritableServers;
97
98
99 private final ClusterType myType;
100
101
102
103
104
105
106
107
108
109 public Cluster(final MongoClientConfiguration config, final ClusterType type) {
110 myConfig = config;
111 myType = type;
112 myChangeSupport = new PropertyChangeSupport(this);
113 myServers = new ConcurrentHashMap<String, Server>();
114 myWritableServers = new CopyOnWriteArrayList<Server>();
115 myNonWritableServers = new CopyOnWriteArrayList<Server>();
116 myListener = new ServerListener();
117 myServerVersionRange = VersionRange.range(Version.parse("0"),
118 Version.parse("0"));
119 }
120
121
122
123
124
125
126
127
128
129 public Server add(final InetSocketAddress address) {
130 final String normalized = ServerNameUtils.normalize(address);
131 Server server = myServers.get(normalized);
132 if (server == null) {
133
134 server = new Server(address);
135
136 synchronized (this) {
137 final Server existing = myServers.putIfAbsent(normalized,
138 server);
139 if (existing != null) {
140 server = existing;
141 }
142 else {
143 myNonWritableServers.add(server);
144 myChangeSupport.firePropertyChange(SERVER_PROP, null,
145 server);
146
147 server.addListener(myListener);
148 }
149 }
150 }
151 return server;
152 }
153
154
155
156
157
158
159
160
161
162
163
164
165
166 public Server add(final String address) {
167 Server server = myServers.get(address);
168 if (server == null) {
169 server = add(ServerNameUtils.parse(address));
170 }
171
172 return server;
173 }
174
175
176
177
178
179
180
181 public void addListener(final PropertyChangeListener listener) {
182 synchronized (this) {
183 myChangeSupport.addPropertyChangeListener(listener);
184 }
185 }
186
187
188
189
190 public void clear() {
191 for (final Server server : myServers.values()) {
192 remove(server);
193 }
194 }
195
196
197
198
199
200
201
202
203
204
205
206 public List<Server> findCandidateServers(final ReadPreference readPreference) {
207 List<Server> results = Collections.emptyList();
208
209 switch (readPreference.getMode()) {
210 case NEAREST:
211 results = findNearestCandidates(readPreference);
212 break;
213 case PRIMARY_ONLY:
214 results = findWritableCandidates(readPreference);
215 break;
216 case PRIMARY_PREFERRED:
217 results = merge(findWritableCandidates(readPreference),
218 findNonWritableCandidates(readPreference));
219 break;
220 case SECONDARY_ONLY:
221 results = findNonWritableCandidates(readPreference);
222 break;
223 case SECONDARY_PREFERRED:
224 results = merge(findNonWritableCandidates(readPreference),
225 findWritableCandidates(readPreference));
226 break;
227 case SERVER:
228 results = findCandidateServer(readPreference);
229 break;
230 }
231
232 return results;
233 }
234
235
236
237
238
239
240
241
242
243
244
245 public List<Server> findServers(final Message message1,
246 final Message message2) {
247 List<Server> servers = Collections.emptyList();
248
249 if (message1 != null) {
250 List<Server> potentialServers = findCandidateServers(message1
251 .getReadPreference());
252 servers = potentialServers;
253
254 if (message2 != null) {
255 servers = new ArrayList<Server>(potentialServers);
256 potentialServers = findCandidateServers(message2
257 .getReadPreference());
258 servers.retainAll(potentialServers);
259 }
260 }
261 return servers;
262 }
263
264
265
266
267
268
269
270
271
272
273
274
275 public Server get(final String address) {
276 return add(address);
277 }
278
279
280
281
282
283
284
285 public List<Server> getNonWritableServers() {
286 return new ArrayList<Server>(myNonWritableServers);
287 }
288
289
290
291
292
293
294
295 public List<Server> getServers() {
296 return new ArrayList<Server>(myServers.values());
297 }
298
299
300
301
302 @Override
303 public VersionRange getServerVersionRange() {
304 return myServerVersionRange;
305 }
306
307
308
309
310
311
312
313
314 @Override
315 public int getSmallestMaxBatchedWriteOperations() {
316 return mySmallestMaxBatchedWriteOperations;
317 }
318
319
320
321
322
323
324
325
326 @Override
327 public long getSmallestMaxBsonObjectSize() {
328 return mySmallestMaxBsonObjectSize;
329 }
330
331
332
333
334
335
336 public ClusterType getType() {
337 return myType;
338 }
339
340
341
342
343
344
345
346 public List<Server> getWritableServers() {
347 return new ArrayList<Server>(myWritableServers);
348 }
349
350
351
352
353
354
355
356 public void remove(final Server server) {
357
358 final Server removed = myServers.remove(server.getCanonicalName());
359 if (removed != null) {
360 removed.removeListener(myListener);
361 myNonWritableServers.remove(removed);
362 myWritableServers.remove(removed);
363
364 updateVersions();
365 }
366 }
367
368
369
370
371
372
373
374 public void removeListener(final PropertyChangeListener listener) {
375 synchronized (this) {
376 myChangeSupport.removePropertyChangeListener(listener);
377 }
378 }
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419 protected final double[] cdf(final List<Server> servers) {
420 Collections.sort(servers, ServerLatencyComparator.COMPARATOR);
421
422
423 final double[] relativeLatency = new double[servers.size()];
424 double sum = 0;
425 double first = Double.NEGATIVE_INFINITY;
426 for (int i = 0; i < relativeLatency.length; ++i) {
427 final Server server = servers.get(i);
428 double latency = server.getAverageLatency();
429
430
431 if (first == Double.NEGATIVE_INFINITY) {
432 first = latency;
433 latency = 1.0D;
434 }
435 else {
436 latency /= first;
437 }
438
439 latency = (1.0D / latency);
440 relativeLatency[i] = latency;
441 sum += latency;
442 }
443
444
445
446 double accum = 0.0D;
447 for (int i = 0; i < relativeLatency.length; ++i) {
448 accum += relativeLatency[i];
449
450 relativeLatency[i] = accum / sum;
451 }
452
453 return relativeLatency;
454 }
455
456
457
458
459
460
461
462
463
464 protected List<Server> findCandidateServer(
465 final ReadPreference readPreference) {
466 final Server server = myServers.get(readPreference.getServer());
467 if ((server != null) && readPreference.matches(server.getTags())) {
468 return Collections.singletonList(server);
469 }
470 return Collections.emptyList();
471 }
472
473
474
475
476
477
478
479
480
481
482
483
484 protected List<Server> findNearestCandidates(
485 final ReadPreference readPreference) {
486 final List<Server> results = new ArrayList<Server>(myServers.size());
487 for (final Server server : myServers.values()) {
488 if (readPreference.matches(server.getTags())) {
489 results.add(server);
490 }
491 }
492
493
494 sort(results);
495
496 return results;
497 }
498
499
500
501
502
503
504
505
506
507
508
509
510
511 protected List<Server> findNonWritableCandidates(
512 final ReadPreference readPreference) {
513 final List<Server> results = new ArrayList<Server>(
514 myNonWritableServers.size());
515 for (final Server server : myNonWritableServers) {
516 if (readPreference.matches(server.getTags())
517 && isRecentEnough(server.getSecondsBehind())) {
518 results.add(server);
519 }
520 }
521
522
523 sort(results);
524
525 return results;
526 }
527
528
529
530
531
532
533
534
535
536
537
538
539
540 protected List<Server> findWritableCandidates(
541 final ReadPreference readPreference) {
542 final List<Server> results = new ArrayList<Server>(
543 myWritableServers.size());
544 for (final Server server : myWritableServers) {
545 if (readPreference.matches(server.getTags())) {
546 results.add(server);
547 }
548 }
549
550
551 sort(results);
552
553 return results;
554 }
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569 protected final void sort(final List<Server> servers) {
570 if (servers.isEmpty() || (servers.size() == 1)) {
571 return;
572 }
573
574
575 final double[] cdf = cdf(servers);
576 final double random = Math.random();
577 int index = Arrays.binarySearch(cdf, random);
578
579
580 if (index < 0) {
581
582 index = Math.abs(index + 1);
583 }
584
585
586
587
588
589
590
591
592
593
594 index = Math.min(cdf.length - 1, index);
595
596
597 Collections.swap(servers, 0, index);
598 }
599
600
601
602
603
604 protected void updateVersions() {
605 Version min = null;
606 Version max = null;
607
608 long smallestMaxBsonObjectSize = Long.MAX_VALUE;
609 int smallestMaxBatchedWriteOperations = Integer.MAX_VALUE;
610
611 for (final Server server : myServers.values()) {
612 min = Version.earlier(min, server.getVersion());
613 max = Version.later(max, server.getVersion());
614
615 smallestMaxBsonObjectSize = Math.min(smallestMaxBsonObjectSize,
616 server.getMaxBsonObjectSize());
617 smallestMaxBatchedWriteOperations = Math.min(
618 smallestMaxBatchedWriteOperations,
619 server.getMaxBatchedWriteOperations());
620 }
621
622 myServerVersionRange = VersionRange.range(min, max);
623 mySmallestMaxBsonObjectSize = smallestMaxBsonObjectSize;
624 mySmallestMaxBatchedWriteOperations = smallestMaxBatchedWriteOperations;
625 }
626
627
628
629
630
631
632
633
634
635 private boolean isRecentEnough(final double secondsBehind) {
636 return ((secondsBehind * 1000) < myConfig.getMaxSecondaryLag());
637 }
638
639
640
641
642
643
644
645
646
647
648 private final List<Server> merge(final List<Server> list1,
649 final List<Server> list2) {
650 List<Server> results;
651 if (list1.isEmpty()) {
652 results = list2;
653 }
654 else if (list2.isEmpty()) {
655 results = list1;
656 }
657 else {
658 results = new ArrayList<Server>(list1.size() + list2.size());
659 results.addAll(list1);
660 results.addAll(list2);
661 }
662 return results;
663 }
664
665
666
667
668
669
670
671
672
673
674 protected final class ServerListener implements PropertyChangeListener {
675 @Override
676 public void propertyChange(final PropertyChangeEvent evt) {
677 final String propertyName = evt.getPropertyName();
678 final Server server = (Server) evt.getSource();
679
680 if (Server.STATE_PROP.equals(propertyName)) {
681
682 final boolean old = !myWritableServers.isEmpty();
683
684 if (Server.State.WRITABLE == evt.getNewValue()) {
685 myWritableServers.addIfAbsent(server);
686 myNonWritableServers.remove(server);
687 }
688 else if (Server.State.READ_ONLY == evt.getNewValue()) {
689 myWritableServers.remove(server);
690 myNonWritableServers.addIfAbsent(server);
691 }
692 else {
693 myWritableServers.remove(server);
694 myNonWritableServers.remove(server);
695 }
696
697 myChangeSupport.firePropertyChange(WRITABLE_PROP, old,
698 !myWritableServers.isEmpty());
699
700 }
701 else if (Server.CANONICAL_NAME_PROP.equals(propertyName)) {
702
703
704
705
706 myServers.remove(evt.getOldValue(), server);
707
708
709 final Server existing = myServers.putIfAbsent(
710 server.getCanonicalName(), server);
711 if (existing != null) {
712
713
714 myNonWritableServers.remove(server);
715 myWritableServers.remove(server);
716 server.removeListener(myListener);
717
718 myChangeSupport.firePropertyChange(SERVER_PROP, server,
719 null);
720 }
721 }
722 else if (Server.VERSION_PROP.equals(propertyName)) {
723
724
725
726 final Version old = (Version) evt.getOldValue();
727
728 if (Version.UNKNOWN.equals(old)
729 || (myServerVersionRange.getUpperBounds()
730 .compareTo(old) <= 0)
731 || (myServerVersionRange.getLowerBounds()
732 .compareTo(old) >= 0)) {
733 updateVersions();
734 }
735 }
736 }
737 }
738 }