softwaremill / softwaremill/macwire

Shared singleton for using with Spark

Aperta
#246 2 commenti 0 reazioni 0 assegnatari Vedi su GitHub

Nessuno ha ancora preso questa issue.

Lingua principale
Scala
Stelle
1.3k
Fork
77
Merge medio
9m
PR unite (30g)
4

Descrizione

When using Spark with external resources like a database, a somehow common pattern is to make the database client shared between tasks so the connection pool is shared. Otherwise, with a large number of tasks/threads, the database connections are exhausted and will lead to issues when scaling.
This rises some complications when using such an object, as it must implement some kind of singleton shared between threads that receive serialized objects.
Any idea on how to do this with MacWire? Any pattern that can be used?

A simple example:

import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession
import org.scalatest.funspec.AnyFunSpec
import org.scalatest.matchers.must.Matchers.{be, convertToAnyMustWrapper}

class ModuleWithSparkSpec extends AnyFunSpec {
  it("runs module with spark") {
    val parallelism = 4
    val module = new Module {
      override lazy val connectionString: String = ""
      override lazy val sparkConf: SparkConf = new SparkConf().setAppName("Test").setMaster(s"local[$parallelism]")
    }

    module.run(parallelism * 3) must be(parallelism * 3) // prints 4 thread ids and 4 different hash codes for 3 times
  }
}

class Runner(val sparkConf: SparkConf, val database: Database) extends Serializable {
  def run(count: Int): Long = {
    val database = this.database
    val sparkConf = this.sparkConf
    SparkSession
      .builder()
      .config(sparkConf)
      .getOrCreate()
      .sparkContext
      .parallelize(0 until count)
      .map { n => database.insert(n) }
      .count()
  }
}

trait Module extends Serializable {
  def run(count: Int): Long = runner.run(count)

  import com.softwaremill.macwire._
  protected lazy val connectionString: String = ""
  protected lazy val sparkConf: SparkConf = new SparkConf().setAppName("").setMaster("")
  protected lazy val database: Database = wire[Database] // this will be serialized and duplicated 4 times
  protected lazy val runner: Runner = wire[Runner]
}

class Database(connectionString: String) extends Serializable with AutoCloseable {
  def insert(n: Int): Unit = {
    println(s"Insert $n on thread id = ${Thread.currentThread().getId}, instance hash code = ${hashCode()}")
  }
  override def close(): Unit = {}
}

So the idea would be to have something instead of wire, or beside, that would make it use a single instance. I was thinking to implement a shared singleton Scope that picks the instance from a concurrent collection, would this be the best way to do it?

  protected lazy val database: Database = sharedSingleton(wire[Database])

Guida per i contributori

Nessuna guida per i contributori indicizzata per questo repository

Come iniziare

  1. Leggi tutta la issue e poi la guida ai contributi del progetto.
  2. Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
  3. Fai un fork del repository e lavora su un branch.
  4. Apri una pull request che faccia riferimento al numero della issue.

Direzione di ricerca

Inizia con l'esempio Spark nell'issue, in particolare Module, Runner e Database, e traccia il modo in cui la serializzazione duplica l'istanza del database tra i task. Il lavoro sarà considerato completato quando sarà definito un pattern o Feature deciso e supportato per condividere il client del database tra i task Spark, e il comportamento e l'utilizzo saranno documentati.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Valutazione

Stack tecnologico
scala
Ambito
distributed-systems
Tipo di issue
Funzionalità
Difficoltà
5/5
Tempo stimato
Più di una settimana
Stato di attività
Ferma
Chiarezza
Da chiarire
Idoneità per principianti
25/100

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.