perri.to: A mashup of things

Canales Waitgroups Y Cancelacion

  2019-10-15


Introducci贸n

Nota: Los materiales de ejemplo de este art铆culo se pueden encontrar aqui. Tener en cuenta que la mayor铆a de las decisiones tomadas a la hora de escribir este c贸digo fueron influenciadas por la necesidad de mostrar como hacer algo as铆 que quizas no sean las mas acertadas y ciertamente este c贸digo no dee ser usado en producci贸n.

Uno de los grandes atractivos de go es la concurrencia (no confundir con paralelismo, algunos buenos ejemplos).

Go nos brinda herramientas varias para trabajar con concurrencias, en este tutorial veremos como hacer uso de canales, waitgroups y contextos con cancelaci贸n que nos permitiran respectivamente comunicarnos entre gorutinas, esperar finalizaci贸n de porciones concurrentes y se帽alar el fin de procesos externos a gorutinas.

Ejercicio de ejemplo

Como ejemplo de c贸digo, haremos una peque帽a herramienta que nos permitir谩 hacer b煤squedas de cadenas de texto arbitrarias en los diferentes sitios regionales de Mercado Libre y comparar resultados del item mas caro entre los varios paises. Es claro que la utilidad es nula pero nos permite construir sobre lo hecho en el post sobre APIs y JSON. (No, Mercado Libre no participa ni condona nada de lo que hacemos, es simplemente un uso de su API abierta que parece estar dentro de lo que permiten los t茅rminos de uso).

Conocimientos Previos.

No es necesario pero si util leer en nuestra publicaci贸n preferida sobre los temas de concurrencia y paralelismo, mas arriba en este mismo art铆culo hay algunos enlaces interesantes que proveen suficiente informaci贸n.

Concurrencia

Para hacer uso de las herramientas antes mencionadas, realizaremos el ejercicio con la siguiente arquitectura.

main()  --> |--> busquedaSitioML(Pais1) -->| --> main()
            |    \_ cotizacion()_/         |
            |--> busquedaSitioML(Pais2) -->|
            |    \_ cotizacion()_/         |
            |--> busquedaSitioML(Pais3) -->|
            |    \_ cotizacion()_/         |
            |--> busquedaSitioML(Pais4) -->|
            |    \_ cotizacion()_/         |

Donde main() ser谩 nuestra rutina principal que invocar谩 a nuestra rutina de b煤squeda en Mercado Libre por cada pais concurrentement y a su vez cada una de estas gorutinas hara su b煤squeda concurrentemente a la b煤squeda de la cotizaci贸n de la moneda de dicho pais contra el Dolar Estado Unidense (como moneda de comparaci贸n entre los varios paises)

No se preocupen si esto no tiene mucho sentido, al final del ejercicio todo ser谩 mas claro.

Precondiciones

Vamos a necesitar como pre-condiciones dos cosas:

  • El texto arbitrario de b煤squeda
  • La lista de sitios oficiales de Mercado Libre

El texto

Obtendremos el texto de los argumentos de la l铆nea de comandos. El paquete incluido os tiene un miembro Args que es un slice de string ([]string) que contiene todos los elementos pertinentes a la llamada del comando que estamos ejecutando, por ejemplo: para ./ejecutable un parametro detras de otro tendr谩 []string{"ejecutable", "un", "parametro", "detras", "de", "otro"} vemos que cada elemento es uno de los argumentos pasados al comando (la separaci贸n es el espacio) tambien podemos observar que el primero es el ejecutable en s铆, por ende para obtener el texto arbitrario tomaremos todos los siguientes con os.Args[1:] que nos devolver谩 un slice con todos menos el primer elemento o si se prefiere, desde el elemento 1 en adelante.

// asignemos los t茅rminos de b煤squeda a un string uniendo los parametros de la linea
// de comando con el string espacio
searchTerms := strings.Join(os.Args[1:], " ")

La lista de sitios

Mercado Libre provee un endpoint para que hagamos un pedido de la lista de sitios, se encuentra en https://api.mercadolibre.com/sites y para obtenerlos usaremos la funci贸n fetchSites() creada para este prop贸sito en el c贸digo de ejemplo, no desarollaremos sobre esta porque es una aplicaci贸n del anterior tutorial basta con mencionar que nos devolver谩 la informaci贸n de los sitios disponibles en forma de un slice del siguiente struct

// mlSite imita la estructura JSON que devuelve la b煤squeda de Sites de Mercado Libre
// en este caso, de un solo site pero la b煤squeda devuelve varios
type mlSite struct {
	DefaultCurrencyID string `json:"default_currency_id"`
	ID                string `json:"id"`
	Name              string `json:"name"`
}

Ya tenemos nuestro texto arbitrario de b煤squeda y nuestra lista de sitios de Mercado Libre.

Vamos a prepararnos para recibir los resultados creando un slice de un struct que creamos para estos efectos:

// siteSearchResult contiene un resultado de b煤squeda, es para uso interno, lo utilizaremos
// para enviar resultados de la gorutina a la rutina principal, contiene todo lo relevante
// que la rutina podria devolver, incluyendo un error por si esta fallara.
type siteSearchResult struct {
	site     mlSite
	price    decimal.Decimal
	priceUSD decimal.Decimal
	ratio    decimal.Decimal
	item     string
	err      error
}

dado que sabemos cuantos resultados habr谩 (tantos como sitios) pre-alocaremos el slice con el tama帽o justo:

	// Hacemos una lista que contendr谩 los resultados de las b煤squedas.
	results := make([]siteSearchResult, 0, len(sites))

Notese que el largo del mismo es 0 y la capacidad es el largo de sites, de este modo podremos asignar nuevos elementos con append pero en ningun momento se re-alocar谩 el Array que hay detras del slice porque la capacidad dada al inicio es justa.

WaitGroups, se帽alizaci贸n entre gorutinas.

Antes de lanzarnos en la creaci贸n de gorutinas, vamos a crear los mecanismos necesarios para saber cuando estas terminaron, hay mas de una forma de lograr esto pero me prece que la mas limpia es utilizando WaitGroups (si quieren investigar alguna otra, por ejemplo, pueden buscar sobre cierre de canales).

Los WaitGroups son una herramienta de sincronizaci贸n entre gorutinas (si bien uno las puede usar libremente en otros contextos, no tienen mucho sentido)

La forma en que estos funcionan es a traves de la creaci贸n de un objeto que es seguro para compartir entre rutinas (no todos los objetos modificables los son ya que varias rutinas podrian competir por una asignaci贸n o una lectura y llevar a condiciones de carrera)

El WaitGroup tiene tres m茅todos a traves de los cuales lograremos sincronizar las gorutinas.

  • Add(int): puede ser invocado tantas veces como sea necesario, los enteros pasados se sumaran a un total, nos dice cuantas gorutinas o ejecuciones esperaremos.
  • Done(): ser谩 invocado por cada rutina o proceso que estamos esperando cuando crea pertinente avisar que lo que se esperaba ya sucedi贸.
  • Wait(): Bloquear谩 la ejecuci贸n hasta que Done haya sido invocado tantas veces como el total de enteros de Add, se puede invocar Add en cualquier momento y Done se puede invocar siempre que el contador de Add sea positivo, si disminuye de 0 causar谩 un panico.

Nota: deber谩 ser un puntero.

	// creamos los WaitGroups para cada una de las go-rutinas que buscar谩.
	wg := &sync.WaitGroup{}
	wg.Add(len(sites))

Luego, dentro de la gorutina llamaremos, al principio de la misma:

func queryForSite(searchCriteria string, site mlSite,
	callerWaiting *sync.WaitGroup, result chan siteSearchResult) {
	// lo primero que haremos es encolar la llamada a Done, del wait group, as铆 cuando
	// esta funci贸n salga, sin importar el resultado se avisar谩 que termin贸 a quien est茅
	// esperando.
	defer callerWaiting.Done()

la llamada es diferida hasta el final de la ejecuci贸n de la funci贸n pra asegurarnos que no indicaremos que estamos listos antes de tiempo, as铆 pues defer se llamar谩 luego del final de la funci贸n.

La funci贸n llamante, a su vez, tendr谩 una llamada bloqueante en el punto en que su ejecuci贸n no pueda continuar sin los resultados de esta gorutina o bien cuando haya terminado todo lo que tiene que hacer y este lista para retornar (si la rutina principal sale el programa terminar谩 y se interrumpir谩 la ejecuci贸n de las gorutinas)

	// esperamos el wait group de todas las gorutinas de b煤squeda, que no terminar谩n hasta
	// que la funcion de procesamiento haya leido su resultado.
	wg.Wait()

Canales, enviando datos entre gorutinas.

Uno de los problemas de la concurrencia es la necesidad de enviar datos entre las gorutinas varias. Podr铆amos utilizar complejas combinaciones de variables y mutexes, como la comunicaci贸n entre threads de varios lenguajes pero Go nos provee una forma nativa de hacer esto, los canales. Los canales son una especie de “tubo” que va entre rutinas y nos permite enviar datos, son tipadas como una variable as铆 que a los fines practicos funciona como una asignaci贸n. Pueden ser con o sin b煤fer, o si se prefiere, tener un b煤fer de 0 a mas lugares. Lo que indicar谩 el tama帽o del b煤fer es cuantos elementos podremos asignar al canal antes de que el mismo bloquee, en este punto asignar al canal bloquear谩 hasta que se lean elementos, en este sentido funciona como una pila FIFO.

Crearemos para nuestro ejemplo un canal de tipo siteSearchResult.

	// creamos un canal, sin buffer, para los resultados.
	resultChannel := make(chan siteSearchResult)

Dada la naturaleza de la gorutina, no podemos retornar valores de la forma tradicional ya que no hay un receptor del otro lado de la asignaci贸n esperando ni una rutina bloqueada. A la hora de retornar utilizaremos el canal que creamos anteriormente y pasamos como par谩metro durante la creaci贸n de la gorutina. El struct que creamos para este proposito y cuyo tipo le asignamos al canal ser谩 el transporte a la hora de hacer esto (notar que tiene campos para error y para 茅xito) de este modo enviaremos uno de esto por el canal en cada punto de salida de la funci贸n, ya sea por error:

	// si fallamos retornamos enseguida.
	if err != nil {
		result <- siteSearchResult{
			site: site,
			err:  err,
		}
		return
	}

como por 茅xito

	result <- siteSearchResult{
		site:     site,
		priceUSD: priceUSD,
		price:    price,
		item:     resultML.Results[0].Title,
		ratio:    currencyRatio,
	}

Dado que no queremos indicar que este proceso se complet贸 hasta que los datos hayan sido consumidos, el canal es sin b煤fer y el envio bloqueante al canal esta en el c贸digo sin ninguna provisi贸n para cancelaci贸n, una vez que este item sea consumido la asignaci贸n dejar谩 de bloquear y la funci贸n prodecer谩 y terminar谩 llamando eventualmente al Done() del wait group, asegurando consistencia.

Alternativamente, como veremos a continuaci贸n, podemos utilizar select para tener alternativas al bloqueo del canal, por ejemplo tiempos m谩ximos de espera o incluso cancelaci贸n externa.

Para consumir los resultados, utilizaremos tambi茅n una gorutina que se comporta de una manera ligeramente distinta:

El consumo de un canal, al igual que la asignaci贸n, es bloqueante. Un patr贸n com煤n de uso para consumir canales que tienen resultados de otras gorutinas es utilizando un loop infinito for {} y consumir el canal dentro de un select que actua como un switch para canales, donde hay varios casos, uno por cada canal con el que queremos interactuar y se dar谩 lugar al primero que se desbloquee, aqu铆 un ejemplo de la funci贸n de consumo:

		for {
			select {
			case r := <-resultChannel:
				if r.err != nil {
					fmt.Printf("Site %q failed %v\n", r.site.Name, r.err)
					break
				}
				results = append(results, r)
			case <-ctx.Done():
				waitResultFetch.Done()
				return
			}
		}

En el ejemplo, el select hace dos operaciones bloqueantes con canales, la primera de ellas intenta asignar lo que salga del canal de resultados resultChannel a r, sabemos que ser谩 nuestro struct de resultado y en base a lo que sea (error o 茅xito) hara lo pertinente. El segundo caso simplemente consume de un canal ctx.Done() (no devuelve nada, simplemente un struct vacio) cuya funci贸n es bloquear hasta que una fuente externa indique que debe seguir, de este modo nunca tendremos un bloqueo infinito ya que podemos indicarle a la rutina que termine incluso i ya nada viene de resultChannel otros usos son, por ejemplo, la funci贸n After(tiempo) de la libreria incluida time que devuelve un canal que al cabo de tiempo devuelve un struct vacio.

Contexto con cancelacion

El context.Context es un elemento algo controversial, hay opiniones algo divididas sobre como usarlo, para que usarlo y si usarlo del todo. Al ser bastante versatil es facil abusar del mismo pero dejaremos estos detalles a criterio del lector ya que es material para un post entero.

En este caso utilizaremos el contexto con Cancelaci贸n, los contextos son objetos que vamos envolviendo con diferentes atributos, tipicamente los modificadores se llaman With* y son funciones que toman un contexto como entrada y devuelven otro que contempla el pasado y ha agregado cosas.

El contexto con cancelaci贸n, que se obtiene pasando un contexto existente (o context.Background si aun no tenemos uno) a context.WithCancel() que nos devuelve el nuevo contexto y una funci贸n que, al invocarla, cancelar谩 el contexto.

	// Hacemos un contexto cancelable para indicar cuando estemos listos
	// para salir de la funci贸n de procesamiento de resultados.
	ctx, done := context.WithCancel(context.Background())

Podemos pasar el elemento ctx y el receptor o receptores pueden invocar ctx.Done() que bloquear谩 hasta que alguien ejecute done() tambien devuelve un canal, de modo tal que se pueda usar en select.

		select {
			case <-ctx.Done():
				//... tareas de cierre ...//
				return
		}

Gorutinas

Finalmente, una peque帽a rese帽a sobre las gorutinas y como las usamos en este ejercicio. Iteraremos sobre todos los sitios e instanciaremos una gorutina por cada uno (esto se logra anteponiendo la palabra clave go al llamado de la funci贸n) y le pasamos el WaitGroup y el chan a cada una.

	// instanciamos una gorutina por cada sitio de Mercado Libre
	for i := range sites {
		go queryForSite(searchTerms, sites[i], wg, resultChannel)
	}

C贸digo terminado

Luego de poner en pr谩ctica todo lo explicado, veremos un ejemplo de como termin贸 nuestro c贸digo en las partes relevantes, pueden encontrar todo el c贸digo con comentarios y explicaciones en el repositorio

El main

func main() {
	// Obtenemos de los argumentos de linea de comandos el criterio de b煤squeda.
	searchTerms := strings.Join(os.Args[1:], " ")

	// obtenemos de mercado libre los sitios internacionales
	sites, err := fetchSites()
	if err != nil {
		log.Fatalf("could not obtain mercado libre sites: %v", err)
	}

	// Hacemos una lista que contendr谩 los resultados de las b煤squedas.
	results := make([]siteSearchResult, 0, len(sites))

	// creamos los WaitGroups para cada una de las go-rutinas que buscar谩.
	wg := &sync.WaitGroup{}
	wg.Add(len(sites))

	// creamos un canal, sin buffer, para los resultados.
	resultChannel := make(chan siteSearchResult)

	// instanciamos una gorutina por cada sitio de Mercado Libre
	for i := range sites {
		go queryForSite(searchTerms, sites[i], wg, resultChannel)
	}

	// creamos un WaitGroup para esperar la gorutina que procesa los resultados.
	waitResultFetch := &sync.WaitGroup{}
	waitResultFetch.Add(1)

	// Hacemos un contexto cancelable para indicar cuando estemos listos
	// para salir de la funci贸n de procesamiento de resultados.
	ctx, done := context.WithCancel(context.Background())

	// invocamos la funci贸n an贸nima de procesamiento de resultados pasando
	// el contexto como par谩metro, notar el shadowing.
	go func(ctx context.Context) {
		for {
			select {
			case r := <-resultChannel:
				if r.err != nil {
					fmt.Printf("Site %q failed %v\n", r.site.Name, r.err)
					break
				}
				results = append(results, r)
			case <-ctx.Done():
				waitResultFetch.Done()
				return
			}
		}
	}(ctx)

	// esperamos el wait group de todas las gorutinas de b煤squeda, que no terminar谩n hasta
	// que la funcion de procesamiento haya leido su resultado.
	wg.Wait()

	// indicamos a la funci贸n de procesamiento que ya no queda nada por procesar
	done()

	// esperamos que la funci贸n de procesamiento termine.
	waitResultFetch.Wait()

	// imprimimos los resultados
	for _, v := range results {
		fmt.Printf("Comprar %q en %q cuesta USD %s (son %s %s a cambio %s):\n",
			searchTerms, v.site.Name, v.priceUSD.StringFixedBank(2), v.site.DefaultCurrencyID, v.price.StringFixedBank(2), v.ratio)
		fmt.Printf("--> Publicado como %q\n", v.item)
	}
}

La funci贸n de b煤squeda de sitios

// queryForSite hara un pedido de b煤squeda y devolver谩 el resultado mas caro para un site
// determinado de Mercado Libre. El resultado se devolver谩 en D贸lares EstadoUnidenses si es
// posible por una cuesti贸n de uniformidad de los resultados (ademas de la moneda de origen)
// esta pensado para ser llamado dentro de una gorutina, concurrentemente con otros sites.
func queryForSite(searchCriteria string, site mlSite,
	callerWaiting *sync.WaitGroup, result chan siteSearchResult) {
	// lo primero que haremos es encolar la llamada a Done, del wait group, as铆 cuando
	// esta funci贸n salga, sin importar el resultado se avisar谩 que termin贸 a quien est茅
	// esperando.
	defer callerWaiting.Done()

	// creamos un wait group para la gorutina que obtendr谩 la cotizaci贸n.
	currencyWait := &sync.WaitGroup{}
	currencyWait.Add(1)
	// como la gorutina es una funci贸n an贸nima dentro de esta, podemos compartir variables
	// para facilitar
	var currencyRatio decimal.Decimal
	var currencyError error

	// llamamos concurrentemente a la funci贸n de b煤squeda de cotizaci贸n, cuando termine
	// lo indicar谩 al wait group.
	go func() {
		defer currencyWait.Done()
		currencyRatio, currencyError = fetchCurrencyRate(site.DefaultCurrencyID)
	}()

	// realizamos la funci贸n principal de esta funci贸n, buscar el item mas caro
	body, err := queryML(searchCriteria, site)
	// si fallamos retornamos enseguida.
	if err != nil {
		result <- siteSearchResult{
			site: site,
			err:  err,
		}
		return
	}

	// leemos el cuerpo de la respuesa
	bodyData, err := ioutil.ReadAll(body)
		// si fallamos retornamos enseguida.
	if err != nil {
		result <- siteSearchResult{
			site: site,
			err:  fmt.Errorf("reading mercado libre response body: %v", err),
		}
		return
	}

	// de-serializamos el cuerpo en un ResultadosML
	resultML := &ResultadosML{}
	err = json.Unmarshal(bodyData, &resultML)
			// si fallamos retornamos enseguida.
	if err != nil {
		result <- siteSearchResult{
			site: site,
			err:  fmt.Errorf("unmarshaling mercado libre response body: %v", err),
		}
		return
	}
			// si no encontramos resultados retornamos enseguida.
	if len(resultML.Results) == 0 {
		result <- siteSearchResult{
			site: site,
			err:  fmt.Errorf("results not found in response"),
		}
		return
	}

	// esperamos a la funci贸n de cotizaci贸n para poder hacer la conversi贸n de moneda.
	currencyWait.Wait()
	// si la funci贸n de cotizaci贸n fall贸, retornaremos enseguida
	if currencyError != nil {
		result <- siteSearchResult{
			site: site,
			err:  fmt.Errorf("getting currency ratio: %v", currencyError),
		}
		return
	}

	// Algunos prints 煤tiles para entender la funci贸n y como se ejecuta.
	//fmt.Println(site.Name)
	//fmt.Println(resultML.Results[0].Title)
	//fmt.Println(resultML.Results[0].Permalink)
	mlResult := resultML.Results[0]
	var price, priceUSD decimal.Decimal
	// si el precio esta en D贸lares EstadoUnidenses originalmente agregaremos la otra
	// cotizaci贸n dividiendo el precio en USD / cotizaci贸n
	// de lo contrario multiplicaremos el precio en moneda de origen por cotizaci贸n para
	// rellenar el precio en USD.
	if mlResult.CurrencyID == usdCurrencyCode {
		priceUSD = mlResult.GetPrice()
		price = priceUSD.Div(currencyRatio)
	} else {
		price = mlResult.GetPrice()
		priceUSD = price.Mul(currencyRatio)
	}

	// enviamos el struct que contiene el resultado por el canal de resultados.
	result <- siteSearchResult{
		site:     site,
		priceUSD: priceUSD,
		price:    price,
		item:     resultML.Results[0].Title,
		ratio:    currencyRatio,
	}
}