Draft: [Interview] Java Concurrency
Interview questions and answers about Java Concurrency
Draft: [Interview] Java Concurrency
This article is an unreviewed draft and may contain incorrect information.
Fundamentals
Thread Lifecycle
NEWRUNNABLEBLOCKEDWAITINGTIMED_WAITINGTERMINATED
Methods
join()- better name would be:
waitForCompletion() - then actual invocation:
parent.waitForCompletion(threadOne)→ parent joinsthreadOneexecution waiting for it to finish
- better name would be:
wait()- give lock back away and wait for
notify()call
- give lock back away and wait for
notify()- notifies the waiting one to wake up
Producer Consumer
Producerproduces >= 1&&Consumerconsumes >=1&&common_buffer
Simple pseudocode:
1
2
3
4
5
6
7
8
// produce
synchronized (this) {
while (buffer == MAX) {
wait()
}
buffer.add()
notifyAll()
}
1
2
3
4
5
6
7
8
// consume
synchronized (this) {
while (buffer == 0) {
wait()
}
buffer.take()
notify()
}
volatile
- used for main-memory visibility guarantees in JVM
- affects happens-before ordering
- makes CPU read/write from common memory e.g. L3 Cache and not its own register
ExecutorService
execute(Runnable)- Schedules aRunnable; does not return a result handle.submit(Runnable)- Schedules aRunnableand returns aFuture<?>to track completion or failure.submit(Callable<T>)- Schedules a value-returning task and returns aFuture<T>.shutdown()- Stops accepting new tasks and lets already submitted tasks finish.shutdownNow()- Attempts to stop active tasks by interrupting workers and returns tasks that never started.awaitTermination(timeout, unit)- Blocks until termination, timeout, or interruption.
SingleThreadExecutor
1
2
3
4
try(ExecutorService service =
Executors.newSingleThreadExecutor()) {
// service.execute(Runnable runnable)
}
FixedThreadPool
1
2
3
4
try(ExecutorService service =
Executors.newFixedThreadPool(int nThreads)) {
// service.execute(Runnable runnable)
}
CachedThreadPool
some kind of “autoscaling” Pool (60s idle → terminate Thread)
1
2
3
4
try(ExecutorService service =
Executors.newCachedThreadPool()) {
// service.execute(Runnable runnable)
}
Scheduled execution
Turn on
1
2
3
4
5
6
7
8
try(ScheduledExecutorService service = Executors.newSingleThreadScheduledExecutor()) {
service.scheduleAtFixedRate(
Runnable task,
long firstDelay,
long period,
TimeUnit timeUnit
);
}
Shut down
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
try {
// true: executor terminated / false: timeout elapsed
boolean succeeded =
service.awaitTermination(
long timeout,
TimeUnit timeUnit
);
if (!succeeded) {
// force shutdown
service.shutdownNow(); // shutdown() would be graceful
}
} catch (InterruptedException e) {
service.shutdownNow();
}
Callable and Future
future.get()- blocking operation, waiting until result arrives
throws ExecutionException, Interrupted Exception
future.get(long timeout, TimeUnit timeUnit)- blocking operation, waiting until timeout
throws TimeoutException(apart fromInterruptedException, ExecutionException)
future.cancel(boolean mayInterruptIfRunning)- attempts to cancel task, may also try to interrupt
future.isCancelled()future.isDone()
1
2
3
4
5
6
7
8
9
Future<Response> response =
service.submit(
new Callable<Response>() {
@Override
public Response call() throws Exception {
return someRequest.invoke();
}
}
);
CompletableFuture
Future<T>answers:- “Has this submitted task finished, and can I obtain its result?”
- Usually you call
get(), which waits.
CompletableFuture<T>answers:- “When this finishes, what should happen next?”
- You attach continuations such as
thenApply(),thenCompose(), andhandle()instead of blocking immediately.
1
2
3
4
5
6
7
8
9
10
11
ExecutorService ioExecutor = Executors.newFixedThreadPool(32);
CompletableFuture<OrderSummary> summaryFuture =
CompletableFuture
.supplyAsync(() -> customerClient.getCustomer(customerId), ioExecutor)
.thenCompose(customer -> CompletableFuture.supplyAsync(() -> orderClient.findRecentOrders(customer.id()), ioExecutor))
.thenApply(orders -> OrderSummary.from(orders))
.exceptionally(ex -> {
log.warn("Unable to build order summary for {}", customerId, ex);
return OrderSummary.unavailable();
});
Synchronized Collections
Manual
- coarse-grained locking (single
Lockfor everything) - limited functionality (no additional methods for locking)
- performance overhead (lock acquisition/release overhead)
1
2
3
4
List<Integer> list =
Collections.synchronizedList(
new ArrayList()
);
BlockingQueue
General behavior:
- taking from empty →
wait - adding to full →
wait
Methods:
put(E e)- put orwaitif fulltake(): E- take orwaitif emptyoffer(E e): boolean-true: element added /false: queue fullpoll(): E?- take from head /null: if emptypeek(): E?- get item without taking /nullif empty
Implementations:
BlockingDeque- double-ended queue (deck)TransferQueue- allowsProducerto directly transfer to waitingConsumerArrayBlockingQueue- bounded queue, backed by fixed arrayLinkedBlockingQueue- (un)bounded queue, backed by linked nodesPriorityBlockingQueue- orders items byComparatorDelayQueue- only expired delay items can be takenSynchronousQueue- zero capacity, only direct transfers
ConcurrentMap
Basic Operations
putIfAbsent(key, value)computeIfAbsent(key, fn)compute()merge(key, value, fn)- inserts a value if absent; otherwise combines old and new values atomically.replace(key, value)- replaces only when the existing value equals oldValue.remove(key, value)- removes only when the current value equals the supplied value.
Implementations:
ConcurrentHashMapConcurrentSkipListMap
CopyOnWriteArrayList
- writers won’t interfere with readers
some like
Gitbranchingadd(element)- appends an element by creating and publishing a new backing-array snapshot.addIfAbsent(element)- adds the element only when it is not already present.set(index, element)- replaces an element by copying and publishing a new array.remove(element)- removes the first matching element by creating a new array.remove(index)- removes the element at the index by creating a new array.get(index)- reads an element without copying the array.contains(element)- checks whether an element exists in the current snapshot.
1
2
3
4
5
6
7
List<Integer> list =
new CopyOnWriteArrayList<>();
// Thread1 can read
sout(list);
// Thread2 can write
list.set(index, value);
Atomic Variables
read-modify-write cycle
count++ is actually:
- Load count value
- Increment loaded value
- Set value to variable
Basic Operations
get()set()compareAndSet(expected, updated)- if (expected) then → updategetAndIncrement()/incrementAndGet()getAndDecrement()/decrementAndGet()
Implementations
AtomicIntegerAtomicBooleanAtomicLongAtomicReferenceLongAdderDoubleAdder
Locks & Synchronizers
CountDownLatch
CountDownLatchis single use, cannot be reused
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
CountDownLatch latch =
new CountDownLatch(int count);
// e.g. for count = 3:
Runnable runnable = () -> {
Thread.sleep(random);
// decrease latch
latch.countDown();
};
new Thread(runnable).start();
new Thread(runnable).start();
new Thread(runnable).start();
// blocks and awaits for all 3 `countDown()` calls
latch.await();
CyclicBarrier
E.g. used for multiple checkpoints during a play when we have to wait for all players to reach given point to proceed
1
2
3
4
5
6
7
CyclicBarrier barrier =
new CyclicBarrier(
int cycle,
Runnable action // what to do when releasing
)
barrier.await();
Exchanger
- synchronization point at which
Threadscan pair and swap elements within concurrent environment - e.g. used if you create pipeline for adjacent steps
1
2
3
4
5
6
7
8
9
Exchanger<String> exchanger =
new Exchanger<>();
// Thread1: calls
String dataFromThread2 =
exchanger.exchange("DataToFlow: Thread1 → Thread2"); // waiting for another `Thread`
// Thread2: calls
String dataFromThread1 =
exchanger.exchange("DataToFlow: Thread2 → Thread1"); // exchange takes place in this step
Condition
Conditioncan be e.g.Queuebeing full →await()- when
Conditionis met, there goes signal →signal()
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
Lock lock = new ReentrantLock();
// Conditions are not programmable, we have to use them properly
// so naming is very important!
Condition bufferNotFull = lock.newCondition();
Condition bufferNotEmpty = lock.newCondition();
// Producer:
lock.lock();
try {
while(buffer.size() == MAX_SIZE) {
bufferNotFull.await(); // wait for condition to be met
}
buffer.offer(item);
bufferNotEmpty.signal(); // signal that condition is now met
} finally {
lock.unlock(); // don't forget to unlock
}
// Consumer:
lock.lock();
try {
while(buffer.size() == 0) {
bufferNotEmpty.await(); // wait until buffer is not empty
}
buffer.poll();
bufferNotFull.signal() // let them know buffer is not full anymore
} finally {
lock.unlock(); // don't forget
}
ReentrantLock
ReentrantLockenables the sameThreadtolock()multiple times without precedingunlock()- It has count how many times lock has been acquired by a given
Thread - Release takes place when counter reaches 0
- Fairness makes
Threadswaiting most time higher in the Priority Queue - Without Fairness result is non-deterministic
Methods
lock()unlock()getHoldCount(): int- count of currentThreadtryLock(): boolean-trueif successfully acquired /falseif not acquiredtryLock(timeout, timeUnit): boolean- similar with timeoutisHeldByCurrentThread(): booleangetQueueLength()newCondition()
Example
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
Lock lock =
new ReentrantLock(boolean fair);
methodA() {
lock.lock();
try {
value++;
methodB();
} finally {
lock.unlock();
}
}
methodB() {
lock.lock();
try {
value--;
} finally {
lock.unlock();
}
}
ReadWriteLock
- Used when resource is read heavy
- Many
Threadscanread - Only one
Threadcanwrite - It contains two separate
Locks- one forReaders, one forWriter - Only one of these
Lockscan be active at a given time Threadsare in Priority Queue and with fairness writers won’t be starvedreadLock()writeLock()
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
ReadWriteLock lock =
new ReentrantReadWriteLock();
writeValue() {
lock.writeLock().lock();
try {
value++;
} finally {
lock.writeLock().unlock();
}
}
getValue() {
lock.readLock().lock();
try {
sout(value);
} finally {
lock.readLock().unlock();
}
}
Semaphore
acquire() / acquire(permits)- acquires permit(s)release() / release(permits)- release permit(s)tryAcquire()tryAcquire(timeout)availablePermits()
1
2
3
4
5
6
Semaphore semaphore =
new Semaphore(int permits, boolean fair);
semaphore.acquire();
// ...
semaphore.release();
StampedLock
TODO optimistic-read alternative to ReadWriteLock
Phaser
TODO more flexible multi-phase alternative to CyclicBarrier
Deadlocks
Deadlock is lock release dependency cycle
jps -l- list processes running Javakill -3 <PID>- kill with Thread Dump
1
2
3
4
5
6
7
8
ThreadMXBean mxBean =
ManagementFactory.getThreadMXBean();
long[] threadIds =
mxBean.findDeadlockedThreads();
ThreadInfo[] threadInfo =
mxBean.getThreadInfo(threadIds);
How to prevent Deadlocks
- Use Timeouts
- Take care of Global Ordering of the Locks
- e.g. always Lock in ascending order (LockA, LockB, LockC)
- Avoid nesting Locks
- Use Thread-Safe alternatives
ForkJoinPool
- similar to
ExecutorService ForkJoincan have subtasks- work stealing
- utilization of multi-core processors
- simplified parallelism
efficient work stealing algorithms
RecursiveTaskRecursiveAction