É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 :
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 (
transformWithStateInPandaset basé sur les lignestransformWithState) 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 :
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 |
|---|---|
Modifier les types de données des types plus restrictifs aux types moins restrictifs. | |
Ajoutez de nouveaux champs au schéma des variables de magasin d'état existantes. | |
Supprimer les champs existants du schéma ou d'une variable d'état du magasin. | |
Réorganiser les champs dans une variable. | |
Ajoutez une nouvelle variable d'état à une application. | |
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
Encoderpour spécifier le schéma de chaque variable. - En Python, vous construisez explicitement un schéma en tant que
StructType.
- En Scala, vous utilisez un
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 :
doublene peut pas être réduit àfloat,longouintfloatne peut pas être réduit àlongouintlongne 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 :
intpeut être promu verslong,floatoudoublelongpeut être promu àfloatoudoublefloatpeut être promu àdoublestringpeut être promu àbytesbytespeut ê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
- Python
// 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
}
}
}
class IntStateProcessor(StatefulProcessor):
def init(self, handle):
# Initial schema with Integer field
state_schema = StructType([
StructField("value1", IntegerType(), True)
])
self.state = handle.getValueState("testState", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
# Convert input value to integer and update state
value = pdf["value"].iloc[0]
self.state.update((int(value),))
# Read current state
current_state = self.state.get()
yield pd.DataFrame({
"id": [key[0]],
"stateValue": [current_state[0]]
})
class LongStateProcessor(StatefulProcessor):
def init(self, handle):
# Later schema with Long field (type widening)
state_schema = StructType([
StructField("value1", LongType(), True)
])
self.state = handle.getValueState("testState", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
# Convert input value to long and update state
value = pdf["value"].iloc[0]
# When reading state written with IntStateProcessor,
# it will be automatically converted to Long
self.state.update((int(value),))
# Read current state
current_state = self.state.get()
yield pd.DataFrame({
"id": [key[0]],
"stateValue": [current_state[0]]
})
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
- Python
// 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
}
}
}
class StateV1Processor(StatefulProcessor):
def init(self, handle):
# Initial schema with a single field
state_schema = StructType([
StructField("value1", IntegerType(), True)
])
self.state = handle.getValueState("testState", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
self.state.update((int(value),))
current_state = self.state.get()
yield pd.DataFrame({
"id": [key[0]],
"stateValue": [current_state[0]]
})
class StateV2Processor(StatefulProcessor):
def init(self, handle):
# Later schema with additional fields
state_schema = StructType([
StructField("value1", IntegerType(), True),
StructField("value2", StringType(), True)
])
self.state = handle.getValueState("testState", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
# Read current state
current_state = self.state.get()
# When reading state written with StateV1(1),
# it will be automatically converted to StateV2(1, None)
value1 = current_state[0]
value2 = current_state[1]
# Now update with both fields populated
self.state.update((int(value), f"metadata-{value}"))
current_state = self.state.get()
yield pd.DataFrame({
"id": [key[0]],
"value1": [current_state[0]],
"value2": [current_state[1]]
})
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
- Python
// 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
}
}
}
class RemoveFieldsOriginalProcessor(StatefulProcessor):
def init(self, handle):
# Initial schema with multiple fields
state_schema = StructType([
StructField("value1", IntegerType(), True),
StructField("value2", StringType(), True)
])
self.state = handle.getValueState("testState", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
self.state.update((int(value), f"metadata-{value}"))
current_state = self.state.get()
yield pd.DataFrame({
"id": [key[0]],
"value1": [current_state[0]],
"value2": [current_state[1]]
})
class RemoveFieldsReducedProcessor(StatefulProcessor):
def init(self, handle):
# Later schema with field removed
state_schema = StructType([
StructField("value1", IntegerType(), True)
])
self.state = handle.getValueState("testState", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
# When reading state written with RemoveFieldsOriginalProcessor(1, "metadata-1"),
# it will be automatically converted to just (1,)
current_state = self.state.get()
value1 = current_state[0]
self.state.update((int(value),))
current_state = self.state.get()
yield pd.DataFrame({
"id": [key[0]],
"value1": [current_state[0]]
})
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
- Python
// 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
}
}
}
class OrderedFieldsProcessor(StatefulProcessor):
def init(self, handle):
# Initial schema with fields in original order
state_schema = StructType([
StructField("value1", IntegerType(), True),
StructField("value2", StringType(), True)
])
self.state = handle.getValueState("testState", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
self.state.update((int(value), f"metadata-{value}"))
current_state = self.state.get()
yield pd.DataFrame({
"id": [key[0]],
"value1": [current_state[0]],
"value2": [current_state[1]]
})
class ReorderedFieldsProcessor(StatefulProcessor):
def init(self, handle):
# Later schema with reordered fields
state_schema = StructType([
StructField("value2", StringType(), True),
StructField("value1", IntegerType(), True)
])
self.state = handle.getValueState("testState", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
# When reading state written with OrderedFieldsProcessor(1, "metadata-1"),
# it will be automatically converted to ("metadata-1", 1)
current_state = self.state.get()
value2 = current_state[0]
value1 = current_state[1]
self.state.update((f"new-metadata-{value}", int(value)))
current_state = self.state.get()
yield pd.DataFrame({
"id": [key[0]],
"value2": [current_state[0]],
"value1": [current_state[1]]
})
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
- Python
// 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
}
}
}
class MultiStateV1Processor(StatefulProcessor):
def init(self, handle):
# Initial schema with a single state variable
state_schema = StructType([
StructField("value1", IntegerType(), True),
StructField("value2", StringType(), True)
])
self.state1 = handle.getValueState("testState1", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
self.state1.update((int(value), f"metadata-{value}"))
current_state = self.state1.get()
yield pd.DataFrame({
"id": [key[0]],
"value1": [current_state[0]],
"value2": [current_state[1]]
})
class MultiStateV2Processor(StatefulProcessor):
def init(self, handle):
# Add a second state variable
state1_schema = StructType([
StructField("value1", IntegerType(), True),
StructField("value2", StringType(), True)
])
state2_schema = StructType([
StructField("value1", StringType(), True),
StructField("value2", IntegerType(), True)
])
self.state1 = handle.getValueState("testState1", state1_schema)
self.state2 = handle.getValueState("testState2", state2_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
self.state1.update((int(value), f"metadata-{value}"))
# Access and update the new state variable
current_state2 = self.state2.get() # Will be None on first run
self.state2.update((f"new-metadata-{value}", int(value)))
current_state1 = self.state1.get()
current_state2 = self.state2.get()
yield pd.DataFrame({
"id": [key[0]],
"state1_value1": [current_state1[0]],
"state1_value2": [current_state1[1]],
"state2_value1": [current_state2[0]],
"state2_value2": [current_state2[1]]
})
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
- Python
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
}
}
}
class MultiStateV2Processor(StatefulProcessor):
def init(self, handle):
# Add a second state variable
state1_schema = StructType([
StructField("value1", IntegerType(), True),
StructField("value2", StringType(), True)
])
state2_schema = StructType([
StructField("value1", StringType(), True),
StructField("value2", IntegerType(), True)
])
self.state1 = handle.getValueState("testState1", state1_schema)
self.state2 = handle.getValueState("testState2", state2_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
self.state1.update((int(value), f"metadata-{value}"))
# Access and update the new state variable
current_state2 = self.state2.get() # Will be None on first run
self.state2.update((f"new-metadata-{value}", int(value)))
current_state1 = self.state1.get()
current_state2 = self.state2.get()
yield pd.DataFrame({
"id": [key[0]],
"state1_value1": [current_state1[0]],
"state1_value2": [current_state1[1]],
"state2_value1": [current_state2[0]],
"state2_value2": [current_state2[1]]
})
class RemoveStateVarProcessor(StatefulProcessor):
def init(self, handle):
# Only use one state variable and delete the other
state_schema = StructType([
StructField("value1", IntegerType(), True),
StructField("value2", StringType(), True)
])
self.state1 = handle.getValueState("testState1", state_schema)
# Delete old state variable that we no longer need
handle.deleteIfExists("testState2")
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
self.state1.update((int(value), f"metadata-{value}"))
current_state = self.state1.get()
yield pd.DataFrame({
"id": [key[0]],
"value1": [current_state[0]],
"value2": [current_state[1]]
})
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
nullpour les champs ajoutés. - Python convertit ces valeurs en
Nonepour 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
0pour les typesintoufalsepour les typesbool.
- Les types de référence renvoient
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
- Python
// 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
}
}
}
class NullDefaultsProcessor(StatefulProcessor):
def init(self, handle):
# Initial schema
state_schema = StructType([
StructField("value1", IntegerType(), True),
StructField("value2", StringType(), True)
])
self.state = handle.getValueState("testState", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
self.state.update((int(value), f"metadata-{value}"))
current_state = self.state.get()
yield pd.DataFrame({
"id": [key[0]],
"value1": [current_state[0]],
"value2": [current_state[1]]
})
class ExpandedNullDefaultsProcessor(StatefulProcessor):
def init(self, handle):
# Evolution: Adding new fields with null/default values
state_schema = StructType([
StructField("value1", IntegerType(), True),
StructField("value2", StringType(), True),
StructField("value3", LongType(), True),
StructField("value4", IntegerType(), True),
StructField("value5", BooleanType(), True)
])
self.state = handle.getValueState("testState", state_schema)
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
for pdf in rows:
value = pdf["value"].iloc[0]
# Reading from state
current_state = self.state.get()
# Showing how null defaults work in Python
# When reading state written with NullDefaultsProcessor state = (1, "metadata-1"),
# it will be automatically converted to (1, "metadata-1", None, None, None)
# In Python, both primitive and reference types will be None
value1 = current_state[0]
value2 = current_state[1]
value3 = current_state[2] # Will be None when evolved from older schema
value4 = current_state[3] # Will be None when evolved from older schema
value5 = current_state[4] # Will be None when evolved from older schema
# Check if value3 is None
if value3 is None:
print("The value3 field is None (default value for evolution)")
value3 = 100 # Set a real value now
# Now update with all fields populated
self.state.update((
value1,
value2,
value3,
value4 if value4 is not None else 42,
value5 if value5 is not None else True
))
current_state = self.state.get()
yield pd.DataFrame({
"id": [key[0]],
"value1": [current_state[0]],
"value2": [current_state[1]],
"value3": [current_state[2]],
"value4": [current_state[3]],
"value5": [current_state[4]]
})
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 |
|
É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 |
|
Examinez attentivement les détails suivants lors du dépannage de l'évolution des schémas pour les variables d'état :
- Certains modèles ne sont pas pris en charge pour l'évolution des schémas. Voir les modèles d'évolution des schémas non pris en charge.
- L'évolution des schémas a toutes les exigences de
transformWithStateet nécessite le format d'encodage Avro. Consultez les exigences. - Vous devez redémarrer une query de streaming pour déployer les modifications de code qui entraînent une évolution des schémas. Voir Quand l'évolution des schémas a-t-elle lieu ?.