Códigos del curso

Snapshots y collects — material extra

12 archivos Java · haz clic en cada uno para desplegarlo

Puedes copiar cada archivo desde aquí, o clonar el repositorio FC_CConcurrente completo:

git clone https://github.com/gilde-valeria/FC_CConcurrente.git

Archivos

Archivos

module-info.java8 líneas

Programas_P5/unam.fc.concurrent.practica5/src/module-info.java

/**
 * 
 */
/**
 * 
 */
module unam.fc.concurrent.practica5 {
}
ExecuteSnapshot.java37 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/ExecuteSnapshot.java

package unam.fc.concurrent.practica5;
// Programa 1: Programa que ejecuta el WFSnapshot, imprime el contenido de los snaps actualizados por update()
//  Utiliza WFSnapshot<T>, el cual utiliza StampedSnap y StampedValue
//
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class ExecuteSnapshot {
    public static void main(String[] args) {
        int capacity = 4;
        int init = -1;
        WFSnapshot<Integer> snapshot = new WFSnapshot<Integer>(capacity,init);
            ExecutorService executor = Executors.newFixedThreadPool(4);
            for(int i = 0; i < 20; i++) {
                final int ntask=i;
                executor.execute(() -> snapshot.update(ntask)); 
            }
            executor.shutdown();
            
            while(!executor.isTerminated()) {};
            
            System.out.println("Snapshot Final");
            StampedSnap<Integer>[] copy = (StampedSnap<Integer>[]) new StampedSnap[capacity];
            copy = snapshot.collect();
            for (int j = 0; j < capacity; j++) {
                System.out.println(" Value: " + copy[j].value + " Owner: " + 
            copy[j].owner + " Stamp: " + copy[j].stamp + " Snap: ");
                
                for (int k = 0; k < capacity; k++) {
                    System.out.println("Thread " + k + " has value: " + copy[j].getSnap(k));
                }
            }
        }
    

}
ExecuteSnapshotQ.java64 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/ExecuteSnapshotQ.java

package unam.fc.concurrent.practica5;
// Programa 2: Programa que ejecuta el WFSnapshotH, imprime el contenido de value del snapshot ssR
// En el ssR se guardan las respuestas y las vistas, las vistas son un scan() del ssI de invocaciones
//  Utiliza WFSnapshotH<T>, el cual utiliza StampedSnapH y StampedValue
//  El WFSnapshotH es diferente al WFSnapshot porque permite guardar todas las operaciones hechas en values y todos los snaps hechos en snap
//  Para ello se modifica StampedSnapH (a diferencia de StampedSnap) y se le agregaron listas para values y snap
import java.util.Random;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class ExecuteSnapshotQ {

    public static void main(String[] args) {
        int capacity = 4;
        String init = null;
        ConcurrentLinkedQueue<Integer> queue = new ConcurrentLinkedQueue<>();// Implementacion linealizable de una cola no acotada
        WFSnapshotH<String> snapshotI = new WFSnapshotH<String>(capacity,init);
        WFSnapshotH<String> snapshotR = new WFSnapshotH<String>(capacity,init);
            ExecutorService executor = Executors.newFixedThreadPool(4);
            
            Random rand = new Random(); //Para crear num randoms
            for(int i = 0; i < 10; i++) {
                int numRand = rand.nextInt(2);//Si es 0 se ejecuta un enq(), si es 1 un deq()
                final int ntask=i;//Los items integers que se agregan a la cola
                executor.execute(new RunnableQ(ntask, queue,numRand,snapshotI,snapshotR)); 
            }
            executor.shutdown();
            
            while(!executor.isTerminated()) {};
            
            // Imprimir los snaps guardados en el update()
            
            /*System.out.println("Snapshot Final in Invs");
            StampedSnapH<String>[] copy = (StampedSnapH<String>[]) new StampedSnapH[capacity];
            copyI = snapshotI.collect();
            for (int j = 0; j < capacity; j++) {
                System.out.println("\n Thread Owner: " + copy[j].owner + " Last Value: " + copy[j].value +  
                        " Last Stamp: " + copy[j].stamp + " Number of Snaps: " + copy[j].snap.size());
                
                  for (int i = 0; i < copy[j].snap.size(); i++) { 
                      System.out.println("The " + i + " snap has: " ); 
                      for (int k = 0; k < capacity; k++) {
                          System.out.println("Thread " + k + " has value: " + copy[j].getSnap(k,i)); }
                  }
                 
            }*/
            
            StampedSnapH<String>[] copy = (StampedSnapH<String>[]) new StampedSnapH[capacity];
            copy = snapshotR.collect();
            
            System.out.println("Snapshot Final --- Values");
            for (int j = 0; j < capacity; j++) {
                System.out.println("\n Thread Owner: " + copy[j].owner +  
                        " Last Stamp: " + copy[j].stamp + " Number of Snaps: " + copy[j].snap.size());
                System.out.println("The values are: "); 
                for (int i = 0; i < copy[j].values.size(); i++) { 
                      
                      System.out.println(copy[j].values.get(i)); }
                  
            }
        }
}
RunnableQ.java57 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/RunnableQ.java

package unam.fc.concurrent.practica5;

import java.util.Arrays;
import java.util.concurrent.Callable;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.stream.Collectors;
import java.util.stream.Stream;

public class RunnableQ<T> implements Runnable{
    private int item;
    ConcurrentLinkedQueue<Integer> queue;
    WFSnapshotH<String> snapshotI;
    WFSnapshotH<String> snapshotR;
    private int num;
    public RunnableQ(int item, ConcurrentLinkedQueue<Integer> queue, 
            int num, WFSnapshotH<String> snapshotI, WFSnapshotH<String> snapshotR){
        this.item = item;
        this.queue = queue;
        this.snapshotI = snapshotI;
        this.snapshotR = snapshotR;
        this.num = num;
    }
    
    @Override
    public void run(){
        try{
            if (this.num == 0) {//Probabilidad de 50/100 de ejecutar enq
                String string = String.format("enq( %s )", this.item);
                snapshotI.update(string);//En el ssI se guardan las invocaciones 
                
                Boolean res = this.queue.add(this.item); // Aplico la operacion de la queue
                
                Object[] result  = (T[]) snapshotI.scan(); // Tomo la foto de las invs escritas
                Stream<Object> stream = Arrays.stream(result); 
                String scan = stream.map( n -> n.toString() )
                .collect( Collectors.joining( " + " ) );
                snapshotR.update(scan + String.format(" |  %s || ", res));// En el ssR se guardan las respuestas

            }
            if (this.num == 1) {//Probabilidad de 50/100 de ejecutar deq
                snapshotI.update("deq()");
                T res = (T) this.queue.poll(); // Regresa null o el valor entero
                
                Object[] result  = (T[]) snapshotI.scan();
                Stream<Object> stream = Arrays.stream(result); 
                String scan = stream.map( n -> n.toString() )
                .collect( Collectors.joining( " + " ) );
                snapshotR.update(scan + String.format(" |  %s || ", res));
            }
            

        }
        catch(Exception e){
        }
    }
}
SimpleSnapshot.java58 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/SimpleSnapshot.java

package unam.fc.concurrent.practica5;

/*
 * SimpleSnapshot.java
 *
 * Created on January 19, 2006, 10:43 AM
 *
 * From "Multiprocessor Synchronization and Concurrent Data Structures",
 * by Maurice Herlihy and Nir Shavit.
 * Copyright 2006 Elsevier Inc. All rights reserved.
 */

import java.util.Arrays;
/**
 *
 * @author Maurice Herlihy
 */

public class SimpleSnapshot<T> implements Snapshot<T> {
  private StampedValue<T>[] a_table;  // array of atomic MRSW registers
  public SimpleSnapshot(int capacity, T init) {
    a_table = (StampedValue<T>[]) new StampedValue[capacity];
    for (int i = 0; i < capacity; i++) {
      a_table[i] = new StampedValue<T>(init);
    }
  }
  public void update(T value) {
    int me = ThreadID.get();
    StampedValue<T> oldValue = a_table[me];
    //oldValue.values.add(value); // Guardar el conjunto de valores hechos
    StampedValue<T> newValue =
        new StampedValue<T>((oldValue.stamp)+1, value);
    a_table[me] = newValue;
  }
  private StampedValue<T>[] collect() {
    StampedValue<T>[] copy = (StampedValue<T>[]) new StampedValue[a_table.length];
    for (int j = 0; j < a_table.length; j++)
      copy[j] = a_table[j];
    return copy;
  }
  public T[] scan() {
    StampedValue<T>[] oldCopy, newCopy;
    oldCopy = collect();
    collect: while (true) {
      newCopy = collect();
      if (! Arrays.equals(oldCopy, newCopy)) {
        oldCopy = newCopy;
        continue collect;
      }
      // clean collect
      T[] result = (T[]) new Object[a_table.length];
      for (int j = 0; j < a_table.length; j++)
        result[j] = newCopy[j].value;
      return result;
    }
  }
}
Snapshot.java23 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/Snapshot.java

package unam.fc.concurrent.practica5;

/*
 * Snapshot.java
 *
 * Created on January 19, 2006, 10:39 AM
 *
 * From "Multiprocessor Synchronization and Concurrent Data Structures",
 * by Maurice Herlihy and Nir Shavit.
 * Copyright 2006 Elsevier Inc. All rights reserved.
 */

/**
 * Interface for snapshot implementations
 * @author Maurice Herlihy
 */
public interface Snapshot<T> {
  
  public void update(T v);
  public T[] scan();
  
}
StampedSnap.java43 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/StampedSnap.java

package unam.fc.concurrent.practica5;

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Future;

/*
 * StampedSnap.java
 *
 * Created on January 19, 2006, 9:03 PM
 *
 * From "Multiprocessor Synchronization and Concurrent Data Structures",
 * by Maurice Herlihy and Nir Shavit.
 * Copyright 2006 Elsevier Inc. All rights reserved.
 */


/**
 * Labeled snapshot value.
 * @param T object type
 * @author Maurice Herlihy
 */
public class StampedSnap<T> extends StampedValue<T> {
    
  public T[] snap; 
  public StampedSnap(T value) {
    super(value);
    snap  = null;
  }
  /**
   * Constructor.
   * @param stamp timestamp
   * @param value object value
   * @param snap 
   */
  public StampedSnap(long stamp, T value, T[] snap) {
    super(stamp, value);
    this.snap = snap;
  }
  public T getSnap(int i) {
      return this.snap[i];
  }
}
StampedSnapH.java51 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/StampedSnapH.java

package unam.fc.concurrent.practica5;


import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Future;

/*
 * StampedSnap.java
 *
 * Created on January 19, 2006, 9:03 PM
 *
 * From "Multiprocessor Synchronization and Concurrent Data Structures",
 * by Maurice Herlihy and Nir Shavit.
 * Copyright 2006 Elsevier Inc. All rights reserved.
 */


/**
 * Labeled snapshot value.
 * @param T object type
 * @author Maurice Herlihy
 */
public class StampedSnapH<T> extends StampedValue<T> {
    
  //public T[] snap; 
    public List<T[]> snap = new ArrayList<T[]>();
    public List<T> values = new ArrayList<T>();
  public StampedSnapH(T value) {
    super(value);
    //snap  = null;
  }
  /**
   * Constructor.
   * @param stamp timestamp
   * @param value object value
   * @param snap 
   */
  public StampedSnapH(long stamp, T value, List<T[]> s, List<T> v) {
    super(stamp, value);
    snap = s;
    values = v;
  }
  public StampedSnapH(long stamp, T value) {
        super(stamp, value);
        
      }
  public T getSnap(int i, int j) {
      return this.snap.get(j)[i];
  }
}
StampedValue.java59 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/StampedValue.java

package unam.fc.concurrent.practica5;
/*
 * StampedValue.java
 *
 * Created on January 19, 2006, 9:49 AM
 *
 * From "Multiprocessor Synchronization and Concurrent Data Structures",
 * by Maurice Herlihy and Nir Shavit.
 * Copyright 2006 Elsevier Inc. All rights reserved.
 */


/**
 * stamps for atomic MRMW register
 * @author Maurice Herlihy
 */
public class StampedValue<T> {
  /**
   * Counter value used for comparison.
   */
  public long stamp;
  /**
   * Thread id of creator
   */
  public int owner;
  /**
   * Register value.
   */
  public T   value;
  /**
   * Least value ever.
   */
  public static final StampedValue MIN_VALUE = new StampedValue(null);
  /**
   * Constructor.
   */
  public StampedValue(long stamp, T value) {
    this.stamp = stamp;
    this.owner = ThreadID.get();
    this.value = value;
  }
  /**
   * Constructor.
   */
  public StampedValue(T value) {
    this.stamp = 0;
    this.value = value;
  }
  public static StampedValue max(StampedValue x, StampedValue y) {
    if (x.stamp > y.stamp) {
      return x;
    } else if (x.owner > y.owner){
      return x;
    } else {
      return y;
    }
  }
}
ThreadID.java40 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/ThreadID.java

package unam.fc.concurrent.practica5;

/*
 * ThreadID.java
 *
 * Created on January 11, 2006, 10:27 PM
 *
 * From "Multiprocessor Synchronization and Concurrent Data Structures",
 * by Maurice Herlihy and Nir Shavit.
 * Copyright 2006 Elsevier Inc. All rights reserved.
 */


/**
 * Assigns unique contiguous ids to threads.
 * @author Maurice Herlihyh
 */
public class ThreadID {
  /**
   * The next thread ID to be assigned
   **/
  private static volatile int nextID = 0;
  /**
   * My thread-local ID.
   **/
  private static ThreadLocalID threadID = new ThreadLocalID();
  public static int get() {
    return threadID.get();
  }
  public static void reset() {
    nextID = 0;
  }
  private static class ThreadLocalID extends ThreadLocal<Integer> {
    protected synchronized Integer initialValue() {
      return nextID++;
    }
  }
}
WFSnapshot.java67 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/WFSnapshot.java

package unam.fc.concurrent.practica5;

/*
 * WFSnapshot.java
 *
 * Created on January 19, 2006, 9:20 PM
 *
 * From "Multiprocessor Synchronization and Concurrent Data Structures",
 * by Maurice Herlihy and Nir Shavit.
 * Copyright 2006 Elsevier Inc. All rights reserved.
 */


/**
 * Wait-free snapshot.
 * @author Maurice Herlihy
 */
public class WFSnapshot<T> implements Snapshot<T> {
  private StampedSnap<T>[] a_table; // array of MRSW atomic registers
  public WFSnapshot(int capacity, T init) {
    a_table = (StampedSnap<T>[]) new StampedSnap[capacity];
    for (int i = 0; i < a_table.length; i++) {
      a_table[i] = new StampedSnap<T>(init);
    }
  }
  public void update(T value) {
    int me = ThreadID.get();
    T[] snap = this.scan();
    StampedSnap<T> oldValue = a_table[me];
    StampedSnap<T> newValue =
      new StampedSnap<T>(oldValue.stamp+1, value, snap);
    a_table[me] = newValue;
  }
  public StampedSnap<T>[] collect() {
    StampedSnap<T>[] copy = (StampedSnap<T>[]) new StampedSnap[a_table.length];
    for (int j = 0; j < a_table.length; j++)
      copy[j] = a_table[j];
    return copy;
  }
  public T[] scan() {
    StampedSnap<T>[] oldCopy;
    StampedSnap<T>[] newCopy;
    boolean[] moved = new boolean[a_table.length];
    oldCopy = collect();
    collect: while (true) {
      newCopy = collect();
      for (int j = 0; j < a_table.length; j++) {
        // did any thread move?
        if (oldCopy[j].stamp != newCopy[j].stamp) {
          if (moved[j]) {       // second move
            return oldCopy[j].snap;
          } else {
            moved[j] = true;
            oldCopy = newCopy;
            continue collect;
          }
        }
      }
      // clean collect
      T[] result = (T[]) new Object[a_table.length];
      for (int j = 0; j < a_table.length; j++)
        result[j] = newCopy[j].value;
      return result;
    }
  }
}
WFSnapshotH.java54 líneas

Programas_P5/unam.fc.concurrent.practica5/src/unam/fc/concurrent/practica5/WFSnapshotH.java

package unam.fc.concurrent.practica5;

public class WFSnapshotH<T> implements Snapshot<T> {
      private StampedSnapH<T>[] a_table; // array of MRSW atomic registers
      public WFSnapshotH(int capacity, T init) {
        a_table = (StampedSnapH<T>[]) new StampedSnapH[capacity];
        for (int i = 0; i < a_table.length; i++) {
          a_table[i] = new StampedSnapH<T>(init);
        }
      }
      public void update(T value) {
        int me = ThreadID.get();
        T[] snap1 = this.scan();
        StampedSnapH<T> oldValue = a_table[me];
        oldValue.snap.add(snap1); // Guardar el conjunto de snaps
        oldValue.values.add(value); // Guardar el conjunto de valores hechos
        StampedSnapH<T> newValue =
          new StampedSnapH<T>(oldValue.stamp+1, value, oldValue.snap, oldValue.values);
        a_table[me] = newValue;
      }
      public StampedSnapH<T>[] collect() {
        StampedSnapH<T>[] copy = (StampedSnapH<T>[]) new StampedSnapH[a_table.length];
        for (int j = 0; j < a_table.length; j++)
          copy[j] = a_table[j];
        return copy;
      }
      public T[] scan() {
        StampedSnapH<T>[] oldCopy;
        StampedSnapH<T>[] newCopy;
        boolean[] moved = new boolean[a_table.length];
        oldCopy = collect();
        collect: while (true) {
          newCopy = collect();
          for (int j = 0; j < a_table.length; j++) {
            // did any thread move?
            if (oldCopy[j].stamp != newCopy[j].stamp) {
              if (moved[j]) {       // second move
                return oldCopy[j].snap.getLast();
              } else {
                moved[j] = true;
                oldCopy = newCopy;
                continue collect;
              }
            }
          }
          // clean collect
          T[] result = (T[]) new Object[a_table.length];
          for (int j = 0; j < a_table.length; j++)
            result[j] = (T) newCopy[j].values; // Regreso el conjunto de valores
          return result;
        }
      }
    }