Aller au contenu principal

Évolution des schémas dans le magasin d'état

Cet article fournit un aperçu de l'évolution des schémas dans le magasin d'état et des exemples de types de modifications de schéma prises en charge.

Qu'est-ce que l'évolution des schémas dans le magasin d'état ?

L'évolution des schémas désigne la capacité d'une application à gérer les changements apportés au schéma des données.

Databricks prend en charge l’évolution des schémas dans le magasin d’état RocksDB pour les applications Structured Streaming qui utilisent transformWithState.

L'évolution des schémas offre une flexibilité pour le développement et la facilité de maintenance. Utilisez l'évolution des schémas pour adapter le modèle de données ou les types de données dans votre magasin d'état sans perdre les informations d'état ou nécessiter un retraitement complet des données historiques.

Exigences

Vous devez définir le format d'encodage du magasin d'état sur Avro pour utiliser l'évolution des schémas. Pour définir cela pour la session en cours, exécutez ce qui suit :

Python
spark.conf.set("spark.sql.streaming.stateStore.encodingFormat", "avro")

L'évolution des schémas est prise en charge uniquement pour les opérations avec état qui utilisent transformWithState ou transformWithStateInPandas. Ces opérateurs ainsi que les APIs et les classes associées ont les exigences suivantes :

  • Disponible dans Databricks Runtime 16.2 et versions ultérieures.
  • Le mode d'accès standard est pris en charge pour Python (transformWithStateInPandas et basé sur les lignes transformWithState) dans Databricks Runtime 16.3 et versions supérieures, et pour Scala (transformWithState) dans Databricks Runtime 17.3 et versions supérieures.
  • RocksDB est le fournisseur de magasin d'état default dans Databricks Runtime 17.3 et les versions ultérieures. Pour les versions de Databricks Runtime antérieures à 17.3, vous devez configurer le fournisseur de magasin d'état RocksDB. Databricks recommande d'activer RocksDB dans le cadre de la configuration compute.

Sur les versions de Databricks Runtime inférieures à 17.3, activez le fournisseur de magasin d'état RocksDB pour la session actuelle en exécutant ce qui suit :

Python
spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

Modèles d'évolution de schémas pris en charge dans le magasin d'état

Databricks prend en charge les modèles d'évolution des schémas suivants pour les Opérations de Structured Streaming avec état.

Modèle

Description

Élargissement de type

Modifier les types de données des types plus restrictifs aux types moins restrictifs.

Ajout de champs

Ajoutez de nouveaux champs au schéma des variables de magasin d'état existantes.

Supprimer des champs

Supprimer les champs existants du schéma ou d'une variable d'état du magasin.

Réorganisation des champs

Réorganiser les champs dans une variable.

Ajout de variables d'état

Ajoutez une nouvelle variable d'état à une application.

Suppression des variables d'état

Supprimez une variable d'état existante d'une application.

Modèle

Description

Élargissement de type

Modifier les types de données des types plus restrictifs aux types moins restrictifs.

Ajout de champs

Ajoutez de nouveaux champs au schéma des variables de magasin d'état existantes.

Supprimer des champs

Supprimer les champs existants du schéma ou d'une variable d'état du magasin.

Réorganisation des champs

Réorganiser les champs dans une variable.

Ajout de variables d'état

Ajoutez une nouvelle variable d'état à une application.

Suppression des variables d'état

Supprimez une variable d'état existante d'une application.

Quand l'évolution des schémas se produit-elle ?

L’évolution des schémas dans le magasin d’état résulte de la mise à jour du code qui définit votre application avec état. Par conséquent, les énoncés suivants s'appliquent :

  • L'évolution des schémas ne se produit pas automatiquement suite aux modifications de schéma dans les données source de la query.
  • L'évolution des schémas ne se produit que lorsqu'une nouvelle version de l'application est déployée. Comme une seule version d’une query de streaming peut s’exécuter simultanément, vous devez redémarrer votre Job de streaming pour faire évoluer le schéma des variables d’état.
  • Votre code définit explicitement toutes les variables d'état et définit le schéma pour toutes les variables d'état.
    • En Scala, vous utilisez un Encoder pour spécifier le schéma de chaque variable.
    • En Python, vous construisez explicitement un schéma en tant que StructType.

Modèles d'évolution des schémas non pris en charge

Les modèles d'évolution des schémas suivants ne sont pas pris en charge :

  • Renommage de champs : Le renommage de champs n'est pas pris en charge, car les champs sont mis en correspondance par nom. Le renommage d'un champ est géré par la suppression du champ et l'ajout d'un nouveau champ. Cette opération n'entraîne pas d'erreur, car la suppression et l'ajout de champs sont autorisés, mais les valeurs du champ d'origine ne sont pas reportées dans le nouveau champ.

  • Possibilité de changement de nom ou de type de clé : Vous ne pouvez pas modifier le nom ou le type des clés dans les variables d'état de carte.

  • Réduction de type Les opérations de réduction de type, également appelées downcasting , ne sont pas prises en charge. Ces opérations peuvent entraîner une perte de données. Voici des exemples d'opérations de réduction de type non prises en charge :

    • double ne peut pas être réduit à float, long ou int
    • float ne peut pas être réduit à long ou int
    • long ne peut pas être restreint à int

Élargissement de type dans le magasin d'état

Vous pouvez étendre les types de données primitifs à des types plus accommodants. Les changements d'élargissement de type suivants sont pris en charge :

  • int peut être promu vers long, float ou double
  • long peut être promu à float ou double
  • float peut être promu à double
  • string peut être promu à bytes
  • bytes peut être promu à string

Les valeurs existantes sont converties en le nouveau type. Par exemple, 12 devient 12.00.

Exemple d'élargissement de type avec transformWithState

Scala
// Initial run with Integer field
case class StateV1(value1: Integer)

class ProcessorV1 extends StatefulProcessor[String, String, String] {
@transient var state: ValueState[StateV1] = _

private val stateV1Encoder = Encoders.product[StateV1]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state = getHandle.getValueState[StateV1](
"testState",
stateV1Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
state.update(StateV1(value.toInt))
value
}
}
}

// Later run with Long field (type widening)
case class StateV2(value1: Long)

class ProcessorV2 extends StatefulProcessor[String, String, String] {
@transient var state: ValueState[StateV2] = _

private val stateV2Encoder = Encoders.product[StateV2]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state = getHandle.getValueState[StateV2](
"testState",
stateV2Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
state.update(StateV2(value.toLong))
value
}
}
}

Ajouter des champs aux valeurs du magasin d'état

Vous pouvez ajouter de nouveaux champs au schéma des valeurs existantes du magasin d'états.

Lors de la lecture de données écrites avec l’ancien schéma, l’encodeur Avro retourne les données des champs ajoutés codées en mode natif en tant que null.

Python interprète toujours ces valeurs comme None. Scala a un comportement default différent selon le type de champ. Databricks recommande d'implémenter une logique pour s'assurer que Scala n'impute pas de valeurs pour les données manquantes. Voir Valeurs par default pour les champs ajoutés à la variable d'état.

Exemples d'ajout de nouveaux champs avec transformWithState

Scala
// Initial run with single field
case class StateV1(value1: Integer)

class ProcessorV1 extends StatefulProcessor[String, String, String] {
@transient var state: ValueState[StateV1] = _

private val stateV1Encoder = Encoders.product[StateV1]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state = getHandle.getValueState[StateV1](
"testState",
stateV1Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
state.update(StateV1(value.toInt))
value
}
}
}

// Later run with additional field
case class StateV2(value1: Integer, value2: String)

class ProcessorV2 extends StatefulProcessor[String, String, String] {
@transient var state: ValueState[StateV2] = _

private val stateV2Encoder = Encoders.product[StateV2]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state = getHandle.getValueState[StateV2](
"testState",
stateV2Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
// When reading state written with StateV1(1),
// it will be automatically converted to StateV2(1, null)
val currentState = state.get()
// Now update with both fields populated
state.update(StateV2(value.toInt, s"metadata-${value}"))
value
}
}
}

Supprimer des champs pour stocker les valeurs d'état

Vous pouvez supprimer des champs du schéma d'une variable existante. Lors de la lecture de données avec l'ancien schéma, les champs présents dans les anciennes données mais pas dans le nouveau schéma sont ignorés.

Exemples de suppression de champs de variables d’état

Scala
// Initial run with multiple fields
case class StateV1(value1: Integer, value2: String)

class ProcessorV1 extends StatefulProcessor[String, String, String] {
@transient var state: ValueState[StateV1] = _

private val stateV1Encoder = Encoders.product[StateV1]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state = getHandle.getValueState[StateV1](
"testState",
stateV1Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
state.update(StateV1(value.toInt, s"metadata-${value}"))
value
}
}
}

// Later run with field removed
case class StateV2(value1: Integer)

class ProcessorV2 extends StatefulProcessor[String, String, String] {
@transient var state: ValueState[StateV2] = _

private val stateV2Encoder = Encoders.product[StateV2]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state = getHandle.getValueState[StateV2](
"testState",
stateV2Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
// When reading state written with StateV1(1, "metadata-1"),
// it will be automatically converted to StateV2(1)
val currentState = state.get()
state.update(StateV2(value.toInt))
value
}
}
}

Réorganiser les champs dans une variable d'état

Vous pouvez réorganiser les champs dans une variable d’état, y compris lorsque vous ajoutez ou supprimez des champs existants. Les champs des variables d’état sont mis en correspondance par nom, et non par position.

Exemples de réorganisation de champs dans une variable d'état

Scala
// Initial run with fields in original order
case class StateV1(value1: Integer, value2: String)

class ProcessorV1 extends StatefulProcessor[String, String, String] {
@transient var state: ValueState[StateV1] = _

private val stateV1Encoder = Encoders.product[StateV1]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state = getHandle.getValueState[StateV1](
"testState",
stateV1Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
state.update(StateV1(value.toInt, s"metadata-${value}"))
value
}
}
}

// Later run with reordered fields
case class StateV2(value2: String, value1: Integer)

class ProcessorV2 extends StatefulProcessor[String, String, String] {
@transient var state: ValueState[StateV2] = _

private val stateV2Encoder = Encoders.product[StateV2]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state = getHandle.getValueState[StateV2](
"testState",
stateV2Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
// When reading state written with StateV1(1, "metadata-1"),
// it will be automatically converted to StateV2("metadata-1", 1)
val currentState = state.get()
state.update(StateV2(s"new-metadata-${value}", value.toInt))
value
}
}
}

Ajouter une variable d'état à une application avec état

Nous pouvons également ajouter des variables d'état entre les exécutions de query.

Remarque : Ce modèle ne nécessite pas d'encodeur Avro et est pris en charge par toutes les applications transformWithState.

Exemple d'ajout d'une variable d'état à une application avec état

Scala
// Initial run with fields in original order
case class StateV1(value1: Integer, value2: String)

class ProcessorV1 extends StatefulProcessor[String, String, String] {
@transient var state1: ValueState[StateV1] = _

private val stateV1Encoder = Encoders.product[StateV1]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state1 = getHandle.getValueState[StateV1](
"testState1",
stateV1Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
state1.update(StateV1(value.toInt, s"metadata-${value}"))
value
}
}
}

case class StateV2(value1: String, value2: Integer)

class ProcessorV2 extends StatefulProcessor[String, String, String] {
@transient var state1: ValueState[StateV1] = _
@transient var state2: ValueState[StateV2] = _

private val stateV1Encoder = Encoders.product[StateV1]
private val stateV2Encoder = Encoders.product[StateV2]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state1 = getHandle.getValueState[StateV1](
"testState1",
stateV1Encoder,
TTLConfig.NONE)
state2 = getHandle.getValueState[StateV2](
"testState2",
stateV2Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
state1.update(StateV1(value.toInt, s"metadata-${value}"))
val currentState2 = state2.get()
state2.update(StateV2(s"new-metadata-${value}", value.toInt))
value
}
}
}

Supprimez une variable d'état d'une application avec état

En plus de supprimer des champs, vous pouvez également supprimer des variables d'état entre les exécutions de query.

Remarque : ce modèle ne nécessite pas d'encodeur Avro et est pris en charge par toutes les transformWithState applications.

Exemple de suppression d'une variable d'état pour une application avec état

Scala
case class StateV1(value1: Integer, value2: String)
case class StateV2(value1: Integer, value2: String)

class ProcessorV1 extends StatefulProcessor[String, String, String] {
@transient var state1: ValueState[StateV1] = _
@transient var state2: ValueState[StateV2] = _

private val stateV1Encoder = Encoders.product[StateV1]
private val stateV2Encoder = Encoders.product[StateV2]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state1 = getHandle.getValueState[StateV1](
"testState1",
stateV1Encoder,
TTLConfig.NONE)
state2 = getHandle.getValueState[StateV2](
"testState2",
stateV2Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
state1.update(StateV1(value.toInt, s"metadata-${value}"))
val currentState2 = state2.get()
state2.update(StateV2(value.toInt, s"new-metadata-${value}"))
value
}
}
}

class ProcessorV2 extends StatefulProcessor[String, String, String] {
@transient var state1: ValueState[StateV1] = _

private val stateV1Encoder = Encoders.product[StateV1]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state1 = getHandle.getValueState[StateV1](
"testState1",
stateV1Encoder,
TTLConfig.NONE)
// delete old state variable that we no longer need
getHandle.deleteIfExists("testState2")
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
state1.update(StateV1(value.toInt, s"metadata-${value}"))
value
}
}
}

default values for fields added to state variable

Lorsque vous ajoutez de nouveaux champs à une variable d'état existante, les variables d'état écrites à l'aide de l'ancien schéma ont le comportement suivant :

  • L'encodeur Avro renvoie une valeur null pour les champs ajoutés.
  • Python convertit ces valeurs en None pour tous les types de données.
  • Le comportement par default de Scala diffère selon le type de données :
    • Les types de référence renvoient null.
    • Les types primitifs renvoient une valeur default, qui diffère selon le type primitif. Les exemples incluent 0 pour les types int ou false pour les types bool.

Il n'existe aucune fonctionnalité intégrée ni métadonnée qui signale le champ comme ayant été ajouté par évolution des schémas. Vous devez implémenter une logique pour gérer les valeurs nulles renvoyées pour les champs qui n'existaient pas dans votre schéma précédent.

Pour Scala, vous pouvez éviter d'imputer des valeurs par default en utilisant Option[<Type>], qui renvoie les valeurs manquantes comme None au lieu d'utiliser le type default.

Vous devez implémenter une logique pour gérer correctement les situations où des valeurs de type None sont retournées en raison de l'évolution des schémas.

Exemple de valeurs default pour les champs ajoutés à une variable d'état

Scala
// Example demonstrating how null defaults work in schema evolution

import org.apache.spark.sql.streaming._
import org.apache.spark.sql.Encoders

// Initial schema that will be evolved
case class StateV1(value1: Integer, value2: String)

class ProcessorV1 extends StatefulProcessor[String, String, String] {
@transient var state: ValueState[StateV1] = _

private val stateV1Encoder = Encoders.product[StateV1]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state = getHandle.getValueState[StateV1](
"testState",
stateV1Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
state.update(StateV1(value.toInt, s"metadata-${value}"))
value
}
}
}

// Evolution: Adding a new field with null/default values
case class StateV2(value1: Integer, value2: String, value3: Long, value4: Option[Long])

class ProcessorV2 extends StatefulProcessor[String, String, String] {
@transient var state: ValueState[StateV2] = _

private val stateV2Encoder = Encoders.product[StateV2]

override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
state = getHandle.getValueState[StateV2](
"testState",
stateV2Encoder,
TTLConfig.NONE)
}

override def handleInputRows(
key: String,
inputRows: Iterator[String],
timerValues: TimerValues): Iterator[String] = {
rows.map { value =>
// Reading from state
val currentState = state.get()

// Showing how null defaults work for different types
// When reading state written with StateV1(1, "metadata-1"),
// it will be automatically converted to StateV2(1, "metadata-1", 0L, None)
println(s"Current state: $currentState")

// For primitive types like Long, the UnsafeRow default for null is 0
val longValue = if (currentState.value3 == 0L) {
println("The value3 field is the default value (0)")
100L // Set a real value now
} else {
currentState.value3
}

// Now update with all fields populated
state.update(StateV2(value.toInt, s"metadata-${value}", longValue))
value
}
}
}

Limitations

Le tableau suivant décrit les limites default pour les modifications d'évolution des schémas :

Description

Default limite

Configuration Spark à remplacer

Évolutions des schémas pour une variable d'état. L'application de plusieurs modifications de schéma lors d'un redémarrage de query est considérée comme une seule évolution des schémas.

16

spark.sql.streaming.stateStore.valueStateSchemaEvolutionThreshold

Évolutions des schémas pour la query en streaming. L'application de plusieurs modifications de schéma lors d'un redémarrage de query est considérée comme une seule évolution des schémas.

128

spark.sql.streaming.stateStore.maxNumStateSchemaFiles

Description

Default limite

Configuration Spark à remplacer

Évolutions des schémas pour une variable d'état. L'application de plusieurs modifications de schéma lors d'un redémarrage de query est considérée comme une seule évolution des schémas.

16

spark.sql.streaming.stateStore.valueStateSchemaEvolutionThreshold

Évolutions des schémas pour la query en streaming. L'application de plusieurs modifications de schéma lors d'un redémarrage de query est considérée comme une seule évolution des schémas.

128

spark.sql.streaming.stateStore.maxNumStateSchemaFiles

Examinez attentivement les détails suivants lors du dépannage de l'évolution des schémas pour les variables d'état :