Concorrenza in C++
std::thread, mutex, lock_guard, atomic, condition_variable, async e future.
Concetti Base
Thread
Flusso di esecuzione indipendente. Più thread condividono la stessa memoria.
Race condition
Due thread accedono alla stessa variabile senza sincronizzazione → risultato imprevedibile.
Mutex
Mutual exclusion: solo un thread alla volta può acquisire il lock. Evita le race conditions.
Deadlock
Due thread si aspettano a vicenda → bloccati per sempre. Da evitare!
g++ -std=c++17 -Wall file.cpp -pthread1. std::thread
#include <iostream>
#include <thread>
#include <vector>
using namespace std;
void saluta(int id) {
cout << "Thread " << id << " dice ciao!\n";
}
void conta(int da, int a, int passo) {
for (int i=da; i<=a; i+=passo)
cout << " " << i << " ";
cout << "\n";
}
int main() {
// Crea un thread
thread t1(saluta, 1);
thread t2(saluta, 2);
// join: aspetta che il thread finisca (OBBLIGATORIO o usa detach)
t1.join();
t2.join();
// Thread con lambda
thread t3([](){ cout << "Lambda thread\n"; });
t3.join();
// Più thread in parallelo
vector<thread> threads;
for (int i=0; i<5; i++) {
threads.emplace_back(conta, i*10, i*10+9, 1);
}
for (auto &t : threads) t.join();
// hardware_concurrency: numero di core disponibili
cout << "CPU cores: " << thread::hardware_concurrency() << "\n";
return 0;
}
2. Mutex e Lock
#include <iostream>
#include <thread>
#include <mutex>
#include <vector>
using namespace std;
// PROBLEMA: race condition senza mutex
int contatore_senza_mutex = 0;
void incrementa_unsafe() {
for (int i=0; i<10000; i++) contatore_senza_mutex++; // NON thread-safe!
}
// SOLUZIONE: mutex
mutex mtx;
int contatore = 0;
void incrementa_safe() {
for (int i=0; i<10000; i++) {
lock_guard<mutex> lock(mtx); // RAII: acquisisce nel costruttore, rilascia nel distruttore
contatore++;
}
}
// unique_lock: più flessibile di lock_guard
void incrementa_unique() {
for (int i=0; i<10000; i++) {
unique_lock<mutex> lock(mtx);
contatore++;
// lock.unlock(); // Puoi rilasciare manualmente
// lock.lock(); // E riaqcuisire
}
}
// Classe thread-safe con lock interno
class ContatoreSicuro {
mutable mutex mtx;
int valore = 0;
public:
void incrementa(int n=1) {
lock_guard<mutex> l(mtx);
valore += n;
}
int get() const {
lock_guard<mutex> l(mtx);
return valore;
}
void reset() {
lock_guard<mutex> l(mtx);
valore = 0;
}
};
int main() {
// Dimostra la race condition
{
vector<thread> ts;
for (int i=0; i<4; i++) ts.emplace_back(incrementa_unsafe);
for (auto &t : ts) t.join();
cout << "Senza mutex: " << contatore_senza_mutex << " (dovrebbe essere 40000)\n";
}
// Con mutex
{
vector<thread> ts;
for (int i=0; i<4; i++) ts.emplace_back(incrementa_safe);
for (auto &t : ts) t.join();
cout << "Con mutex: " << contatore << " (esattamente 40000)\n";
}
// Classe thread-safe
ContatoreSicuro cs;
{
vector<thread> ts;
for (int i=0; i<4; i++) ts.emplace_back([&cs](){ for(int j=0;j<10000;j++) cs.incrementa(); });
for (auto &t : ts) t.join();
}
cout << "ContatoreSicuro: " << cs.get() << "\n"; // Sempre 40000
return 0;
}
3. std::atomic
Per tipi semplici (int, bool, pointer), std::atomic fornisce operazioni thread-safe senza mutex (più veloce).
#include <iostream>
#include <thread>
#include <atomic>
#include <vector>
using namespace std;
atomic<int> contatore_atomico{0};
atomic<bool> stop_flag{false};
void worker(int n) {
for (int i=0; i<n; i++) {
contatore_atomico++; // Operazione atomica
// contatore_atomico.fetch_add(1); // Equivalente esplicito
}
}
void timer_thread() {
this_thread::sleep_for(100ms);
stop_flag = true; // Scrittura atomica
}
void loop_until_stop() {
int count = 0;
while (!stop_flag) { // Lettura atomica
count++;
}
cout << "Loop eseguito " << count << " volte\n";
}
int main() {
vector<thread> ts;
for (int i=0; i<4; i++) ts.emplace_back(worker, 10000);
for (auto &t : ts) t.join();
cout << "Atomico: " << contatore_atomico << "\n"; // Sempre 40000
// Operazioni avanzate
atomic<int> val{10};
int vecchio = val.exchange(20); // Scambia e restituisce il vecchio
cout << "exchange: " << vecchio << " → " << val << "\n";
// compare_and_swap (CAS) — fondamentale per strutture lock-free
int expected = 20;
bool successo = val.compare_exchange_strong(expected, 30);
cout << "CAS " << (successo?"riuscito":"fallito") << " → " << val << "\n";
return 0;
}
4. condition_variable e async/future
#include <iostream>
#include <thread>
#include <mutex>
#include <condition_variable>
#include <queue>
#include <future>
#include <chrono>
using namespace std;
// ===== Producer-Consumer con condition_variable =====
queue<int> buffer;
mutex mtx;
condition_variable cv;
const int MAX_BUFFER = 10;
bool finito = false;
void producer(int n) {
for (int i=0; i<n; i++) {
unique_lock<mutex> lock(mtx);
cv.wait(lock, []{ return buffer.size() < MAX_BUFFER; }); // Aspetta se pieno
buffer.push(i);
cout << "Prodotto: " << i << "\n";
cv.notify_all(); // Notifica i consumer
}
{
lock_guard<mutex> lock(mtx);
finito = true;
}
cv.notify_all();
}
void consumer(int id) {
while (true) {
unique_lock<mutex> lock(mtx);
cv.wait(lock, []{ return !buffer.empty() || finito; });
if (buffer.empty()) break; // Finito e buffer vuoto
int val = buffer.front(); buffer.pop();
lock.unlock();
cv.notify_all();
cout << "Consumer " << id << " consuma: " << val << "\n";
}
}
// ===== std::async — esecuzione asincrona =====
int calcola_pesante(int n) {
this_thread::sleep_for(chrono::milliseconds(100)); // Simula calcolo lento
int sum = 0;
for (int i=1; i<=n; i++) sum += i;
return sum;
}
int main() {
// async: lancia in background, ottieni il risultato con future::get()
future<int> risultato = async(launch::async, calcola_pesante, 100);
// Fai altro mentre calcola...
cout << "Calcolo in corso...\n";
this_thread::sleep_for(chrono::milliseconds(50));
cout << "Ancora in attesa...\n";
// get() blocca fino a quando il risultato è pronto
cout << "Risultato: " << risultato.get() << "\n";
// Lancia più calcoli in parallelo
vector<future<int>> futures;
for (int i=1; i<=4; i++)
futures.push_back(async(launch::async, calcola_pesante, i*1000));
int totale = 0;
for (auto &f : futures) totale += f.get();
cout << "Totale parallelo: " << totale << "\n";
// promise/future: comunicazione tra thread
promise<int> prom;
future<int> fut = prom.get_future();
thread t([&prom](){
this_thread::sleep_for(50ms);
prom.set_value(42); // Invia il valore al future
});
cout << "Attendo il valore...\n";
cout << "Ricevuto: " << fut.get() << "\n";
t.join();
return 0;
}
Dividi un array di 1.000.000 elementi in N parti (dove N = numero di core). Lancia un thread per ogni parte che calcola la somma parziale. Usa std::atomic<long long> per accumulare. Confronta il tempo con la versione single-thread.
#include <iostream>
#include <thread>
#include <atomic>
#include <vector>
#include <chrono>
#include <numeric>
using namespace std;
int main() {
const int N = 10000000;
vector<int> arr(N);
iota(arr.begin(), arr.end(), 1); // 1..N
auto ora = []{ return chrono::high_resolution_clock::now(); };
auto ms = [](auto t1, auto t2){
return chrono::duration_cast<chrono::microseconds>(t2-t1).count() / 1000.0;
};
// Single thread
auto t1 = ora();
long long sum_seq = accumulate(arr.begin(), arr.end(), 0LL);
auto t2 = ora();
cout << "Sequenziale: " << sum_seq << " (" << ms(t1,t2) << " ms)\n";
// Multi-thread
int n_thread = thread::hardware_concurrency();
atomic<long long> sum_par{0};
vector<thread> ts;
int chunk = N / n_thread;
auto t3 = ora();
for (int i=0; i<n_thread; i++) {
int da=i*chunk, a=(i==n_thread-1)?N:(i+1)*chunk;
ts.emplace_back([&arr, &sum_par, da, a](){
long long s=0;
for(int j=da;j<a;j++) s+=arr[j];
sum_par += s;
});
}
for (auto &t : ts) t.join();
auto t4 = ora();
cout << "Parallelo (" << n_thread << " thread): " << sum_par
<< " (" << ms(t3,t4) << " ms)\n";
cout << "Speedup: " << ms(t1,t2)/ms(t3,t4) << "x\n";
return 0;
}Implementa una classe ThreadSafeQueue<T> con le operazioni push, pop (bloccante), try_pop (non bloccante), empty, e una variante con timeout. Testa con produttore e consumatore in thread separati.
#include <iostream>
#include <queue>
#include <mutex>
#include <condition_variable>
#include <thread>
#include <optional>
#include <chrono>
using namespace std;
template<typename T>
class ThreadSafeQueue {
queue<T> q;
mutable mutex mtx;
condition_variable cv;
bool closed = false;
public:
void push(T val) {
{
lock_guard<mutex> l(mtx);
q.push(move(val));
}
cv.notify_one();
}
// Bloccante
optional<T> pop() {
unique_lock<mutex> l(mtx);
cv.wait(l, [this]{ return !q.empty() || closed; });
if (q.empty()) return nullopt;
T val = move(q.front()); q.pop();
return val;
}
// Non bloccante
optional<T> try_pop() {
lock_guard<mutex> l(mtx);
if (q.empty()) return nullopt;
T val = move(q.front()); q.pop();
return val;
}
void close() {
{ lock_guard<mutex> l(mtx); closed=true; }
cv.notify_all();
}
bool empty() const {
lock_guard<mutex> l(mtx);
return q.empty();
}
};
int main() {
ThreadSafeQueue<int> queue;
thread prod([&queue](){
for(int i=0;i<10;i++){
queue.push(i);
this_thread::sleep_for(chrono::milliseconds(10));
}
queue.close();
});
thread cons([&queue](){
while(auto val = queue.pop()) {
cout << "Ricevuto: " << *val << "\n";
}
cout << "Queue chiusa\n";
});
prod.join();
cons.join();
return 0;
}