Menu Docs
Página inicial do Docs
/ /

sp.processor.modify() (método mongosh)

sp.processor.modify()

Modifica um Stream Processor nomeado no Stream Processing Workspaceatual.

Esse método é suportado em Atlas Stream Processing Workspaces.

O método sp.processor.modify() tem a seguinte sintaxe:

sp.processor.modify({
pipeline: [
<pipeline>
],
name: <name>,
dlq: {
connectionName: <connectionName>,
db: <db>,
coll: <coll>
},
resumeFromCheckpoint: <resumeFromCheckpoint>,
tier: <tier>
})

sp.processor.modify() usa os seguintes campos:

Campo
Tipo
necessidade
Descrição

pipeline

array

Opcional

Array de estágios de agregação para aplicar aos seus dados de streaming onde o último estágio deve ser um estágio de sink. Para saber mais, consulte Agregação de Stream Processing.

name

string

Opcional

Novo nome para o processador de stream.

dlq

objeto

Opcional

Objeto que define uma fila de mensagens não entregues (DLQ) para seu processador de fluxo. Para remover uma fila de mensagens não entregues (DLQ) existente, passe um objeto vazio ({}).

dlq.connectionName

string

Condicional

Rótulo que identifica uma conexão em seu registro de conexão. Esta conexão deve fazer referência a um cluster do MongoDB Atlas. Obrigatório ao definir uma fila de mensagens não entregues (DLQ).

dlq.db

string

Condicional

Nome de um banco de dados MongoDB Atlas no cluster especificado em dlq.connectionName. Obrigatório ao configurar uma fila de mensagens não entregues (DLQ).

dlq.coll

string

Condicional

Nome de uma coleção no banco de dados especificado em dlq.db. Obrigatório ao configurar uma fila de mensagens não entregues (DLQ).

resumeFromCheckpoint

booleano

Opcional

Sinalizador que especifica se o processador de fluxo modificado é retomado do último checkpoint. Por padrão, este campo é true. Quando definido como false, o processador retém apenas estatísticas de resumo.

tier

string

Opcional

O tier a ser atribuído ao processador de stream. Se você não declarar esta opção, o processador manterá seu tier atual. Deve ser um dos seguintes valores:

  • SP2

  • SP5

  • SP10

  • SP30

  • SP50

Para saber mais, consulte Níveis.

O processador de fluxo deve estar em um estado STOPPED antes de invocar este método. O argumento pipeline substitui todo o pipeline existente do processador, incluindo estágios que você não altera.

Por padrão, o processador modificado é retomado do último checkpoint. Se você definir resumeFromCheckpoint como false, o processador modificado reterá apenas estatísticas de resumo. Quando você modifica um processador com janelas abertas, o Atlas Stream Processing recalcula essas janelas no pipeline atualizado.

Para limitações que se aplicam ao modificar processadores de fluxo, consulte Modificar um processador de fluxo.

Para executar sp.processor.modify(), você deve ter a função atlasAdmin.

O exemplo altera um processador de fluxo interrompido chamado solarDemo para adicionar um estágio $match, renomeá-lo, atualizar seu tier e configurar uma fila de mensagens não entregues (DLQ):

sp.solarDemo.modify({
pipeline: [
{ $source: { connectionName: "sample_stream_solar" }},
{ $match: { device_id: "device_0" }},
{ $merge: { into: {
connectionName: "cluster0",
db: "testout",
coll: "testout2"
}}}
],
name: "solarDemoRenamed",
dlq: {
connectionName: "cluster0",
db: "testout",
coll: "dlq"
},
resumeFromCheckpoint: true,
tier: "SP10"
})
{ ok: 1 }

Inicie o processador renomeado e, em seguida, execute sp.listStreamProcessors() para verificar as alterações de nome, tier e fila de mensagens não entregues (DLQ):

sp.solarDemoRenamed.start()
sp.listStreamProcessors()
[
{
id: '6a39b08e6d9040e1cef8e31f',
name: 'solarDemoRenamed',
lastModified: ISODate('2026-06-22T22:00:46.858Z'),
state: 'STARTED',
tier: 'SP10',
errorMsg: '',
workers: [ 'worker-5f4c5bbc9d-7hg2q' ],
pipeline: [
{ '$source': { connectionName: 'sample_stream_solar' } },
{ '$match': { device_id: 'device_0' } },
{
'$merge': {
into: {
connectionName: 'cluster0',
db: 'testout',
coll: 'testout2'
}
}
}
],
lastStateChange: ISODate('2026-06-22T22:01:16.835Z'),
dlq: {
connectionName: 'cluster0',
db: 'testout',
coll: 'dlq'
}
}
]

Execute sp.processor.sample() para verificar a alteração do pipeline:

sp.solarDemoRenamed.sample()
{
device_id: 'device_0',
group_id: 9,
timestamp: '2026-06-22T22:01:25.828+00:00',
max_watts: 450,
event_type: 0,
obs: {
watts: 122,
temp: 18
}
}
{
device_id: 'device_0',
group_id: 3,
timestamp: '2026-06-22T22:01:26.828+00:00',
max_watts: 450,
event_type: 0,
obs: {
watts: 377,
temp: 7
}
}

Voltar

start.processor.drop

Nesta página