September 3

Пишем ThreadSafe код без Mutex или что такое этот ваш Lock-free?

В прошлых заметках мы прошлись по синтетической задаче TransferMoney вдоль и поперек. Сошлись на том что код с Mutex - самая простая и понятная реализация. Но есть ли у нее альтернативы? Мы попытались написать код на атомиках, он оказался сложнее и имел баги. В этой заметке я попробую написать код без найденных недостатков.

Ахиллесова пята нашего первого "блина комом" - наличие двух атомиков. Достичь атомарности и предсказуемости над двумя участками памяти в лоб невозможно, наше текущее железо к такому просто не приспособлено. Но потребность писать код подобный нашему синтетическому примеру была есть и будет. И раз уж железяка нам не может помочь - будем что-то придумывать на уровне алгоритма самой программы.

Multi-Word Compare-And-Swap

В 2002м году вышла статья A Practical Multi-Word Compare-and-Swap Operation by Timothy L. Harris, Keir Fraser and Ian A. Pratt. В ней был предложена реализация алгоритма CASN - гарантированное атомарное изменение N atomic переменных. Я не буду ее пересказывать, авторы постарались на славу и не только предложили алгоритм но и доказали что он рабочий и корректный (а это в теме concurrency самое главное).

Главные мысли которые нам нужны чтобы приступить к кодированию:

  • Вводим примитив "дескриптор" - он содержит внутри себя операции которые нужно применить атомарно.
  • Вместо мьютексов мы рядом с нашими атомиками хранящими ценные данные храним еще и ссылку на дескриптор - это сигнал программе что над этим кусочком данных кто-то уже начал работу.

Вместе эти два фактора дают нам возможность внедрить логику "доталкивания". Пример: если поток №1 начал операцию и его внезапно усыпил ОС то поток №2 которому нужна та же ячейка памяти для своей операции вместо того чтобы блокироваться и засыпать поток №2 пытается применить ту самую операцию. Вспоминаем определение lock-free - код гарантирует что в случае конкурентного исполнения кто-то обязательно достигнет прогресса.

Похоже что авторам удалось придумать то что мы так долго искали. Давайте попробуем приземлить идеи на практику.

Show me your code

TLDR: Весь код в формате Go Playground доступен здесь. Переходите и экспериментируйте. Перейдем к рассмотрению кода, он будет состоять из нескольких примитивов:

  • Аккаунт - место где мы храним денежку.
  • Транзакция - абстракция хранящая наши операции, гарантирующиая детерминированное исполнение.
  • Операция - примитив хранящий данные о том в каком состоянии данные сейчас, и какой инкремент нужно произвести.

Начнем со структуры account

type Account struct {
	balance atomic.Pointer[balance]
}

type balance struct {
    amount int64
    tx *tx
}

func (b *balance) isInFlight() bool { return b.tx != nil }
func NewAccount(amount int64) *Account {
	a := &Account{}
	a.balance.Store(&balance{amount: amount})
	return a
}

func (a *Account) Balance() int64 {
	return a.readBalance().amount
}

// readBalance возвращает реальное значение ячейки,
// помогая завершить чужую транзакцию если ячейка захвачена.
func (a *Account) readBalance() *balance {
	for {
		b := a.balance.Load()
		if !b.isInFlight() {
			return b
		}
		b.tx.commit()
	}
}

// Transfer — lock-free перевод через Multi-word CAS.
func (from *Account) Transfer(to *Account, amount int64) {
	for {
		t := &tx{}
		t.firstOp = newOperation(from, t, -amount)
		t.secondOp = newOperation(to, t, +amount)

		// Фиксируем порядок операций по адресу, чтобы "доталкивание" работало корректно:
		// без порядка A помогает B, B помогает A → бесконечная рекурсия.
		if uintptr(unsafe.Pointer(from)) > uintptr(unsafe.Pointer(to)) {
			t.firstOp, t.secondOp = t.secondOp, t.firstOp
		}

		if t.commit(); t.succeeded() {
			return
		}
	}
}

Что мы видим:

  • Бесконечный цикл внутри функции Transfer, классика non-blocking алгоритмов. Пробуем сделать что-то полезнон пока не достигнем успеха.
  • Чтение баланса - не просто возвращает нам цифру - под капотом идет проверка на наличии inFlight операции, и если есть - мы ее "доталкиваем" как писали в статье.

Перейдем теперь на уровень ниже - реализация примитива транзакция:

type txStatus int32

const (
	txUnresolved txStatus = iota
	txSucceeded
	txFailed
)

// tx — дескриптор атомарной операции над двумя счетами.
type tx struct {
	done     atomic.Int32
	firstOp  operation
	secondOp operation
}

func (t *tx) status() txStatus { return txStatus(t.done.Load()) }
func (t *tx) resolved() bool   { return t.status() != txUnresolved }
func (t *tx) succeeded() bool  { return t.status() == txSucceeded }
func (t *tx) tryResolve(ok bool) {
	if ok {
		t.done.CompareAndSwap(int32(txUnresolved), int32(txSucceeded))
	} else {
		t.done.CompareAndSwap(int32(txUnresolved), int32(txFailed))
	}
}

func (t *tx) commit() {
	if !t.resolved() {
		t.tryResolve(t.prepare())
	}
	if t.succeeded() {
		t.applyProgress()
	} else {
		t.applyRollback()
	}
}

func (t *tx) prepare() bool {
	return t.firstOp.prepare() && t.secondOp.prepare()
}

func (t *tx) applyProgress() {
	t.firstOp.tryFinalize()
	t.secondOp.tryFinalize()
}

func (t *tx) applyRollback() {
	t.firstOp.tryRestore()
	t.secondOp.tryRestore()
}

Что мы видим:

  • Транзакция может быт в статусе - unresolved, success, failure
  • Транзация состоит из двух фаз - подготовка (prepare) и фиксация (commit)
  • если по каким то причинам не удалось успешно подготовить все вложенные в транзакцию операции то мы падаем и пробуем заново (см код выше)

Перейдем к заключительной детали паззла - примитив операция:

 type operation struct {
	acc    *Account
	before *balance
	after  *balance
}

func newOperation(acct *Account, t *tx, delta int64) operation {
	current := acct.readBalance()
	return operation{
		acc:    acct,
		before: current,
		after:  &balance{amount: current.amount + delta, tx: t},
	}
}

func (op *operation) tx() *tx                     { return op.after.tx }
func (op *operation) current() *balance           { return op.acc.balance.Load() }
func (op *operation) isClaimed(cur *balance) bool { return cur == op.after }
func (op *operation) isStale(cur *balance) bool   { return cur != op.before }
func (op *operation) tryClaim() bool {
	return op.acc.balance.CompareAndSwap(op.before, op.after)
}

func (op *operation) tryFinalize() {
	op.acc.balance.CompareAndSwap(
	    op.after, 
	    &balance{amount: op.after.amount},
	)
}

func (op *operation) tryRestore() {
	op.acc.balance.CompareAndSwap(op.after, op.before)
}

func (op *operation) prepare() bool {
	for {
		if op.tx().resolved() {
			return op.tx().succeeded()
		}
		cur := op.current()
		switch {
		case op.isClaimed(cur):
			return true
		case cur.isInFlight():
			cur.tx.commit()
			continue
		case op.isStale(cur):
			return false
		default:
			op.tryClaim()
		}
	}
}

Что мы видим:

  • Бесконечный цикл, крутимся пока либо сами не допушим транзакцию либо нам не поможет сосед. Ну или мы настолько отстали что нужно падать с false чтобы управляющий код сделал ретрай.

Опытные ребята наверняка сейчас поймают эффект легкого дежавю. Код что-то напоминает. Действительно, то что у нас получилось это оптимистичный 2 phase commit. Еще раз убеждаемся что Concurrency и Distributed Systems идут рука об руку.

Завершаем разбор исходного кода, функция main:

func main() {
    // простой пример, без concurrency
	a := NewAccount(1000)
	b := NewAccount(0)
	fmt.Printf("before: a=%d, b=%d\n", a.Balance(), b.Balance())
	a.Transfer(b, 300)
	fmt.Printf("after:  a=%d, b=%d\n", a.Balance(), b.Balance())

	// много параллельных операций, убедимся что 
	// не теряем и не делаем леньги из воздуха :)
	const (
		numAccounts   = 10
		initialFunds  = 1000
		numGoroutines = 50
		numTransfers  = 200
	)

	accounts := make([]*Account, numAccounts)
	for i := range accounts {
		accounts[i] = NewAccount(initialFunds)
	}
	wantTotal := int64(numAccounts * initialFunds)

	var wg sync.WaitGroup
	for range numGoroutines {
		wg.Add(1)
		go func() {
			defer wg.Done()
			rng := rand.New(rand.NewSource(rand.Int63()))
			for range numTransfers {
				from := accounts[rng.Intn(numAccounts)]
				to := accounts[rng.Intn(numAccounts)]
				if from == to {
					continue
				}
				from.Transfer(to, int64(rng.Intn(10)+1))
			}
		}()
	}
	wg.Wait()

	var total int64
	for _, acc := range accounts {
		total += acc.Balance()
	}
	fmt.Printf("conservation check: want=%d got=%d ok=%v\n", wantTotal, total, total == wantTotal)
}

Выводы или "Зачем так сложно?"

Нам удалось достичь поставленной цели - реализовать алгоритм без мьютексов в котором разделяемая память представлена сложнее чем единственный atomic.

Если вы дочитали до этого места у вас наверняка возник вопрос вынесенный в заголовок. И он абсолютно справедливый и здравый. Скорее всего в продакшне напрямую такой никто из нас никогда не писал и даже не видел. Скажу больше - именно в таком виде встретить реализацию структуры данных и алгоритма практически невозможно - слишком наивно и дорого по перформансу.

Но так или иначе идеи которые мы рассмотрели это фундамент на основе которого построены например:

  • Software Transactional Memory подход который очень классно себя показал в функциональных языках программирования. Про него сделаю отдельный пост в канале или заметку.
  • Lock-free структуры данных (идеи MCAS нашли отклик в реализации дереьев поиска, а идея "доталкивания" присутствует дефакто стандартной lock-free реализации очереди - Michael-Scott lock-free queue.

На этом всё, надеюсь вам было интересно и не скучно. В следующих постах мы будем уходить подальше от академических изысков и будем искать ответы на вопросы вида:

  • А что по перформансу? Lock-Free уделывает мьютексы?
  • Почему мы все еще пользуемся мьютексами если можно без них?

Спасибо, что дочитали, до встречи в новых постах!