Definição
A expressão $meta retorna um objeto que contém todos os metadados de streaming de um documento. Você pode expor esses dados para todo o fluxo ou para um dos seguintes estágios de agregação do Atlas Stream Processing:
Uma expressão $meta tem o seguinte formato de protótipo:
{ "$meta": <string> }
"source": { "type": "<source-type>", "ts": { "$date": "<datetime>" }, "topic": "<string>", "partition": <int>, "offset": <int>, "key": "<kafka-key>", "headers": [ { "k": "<header-key>", "v": "<header-value>" } ], "operationType": "<db-operation>", "ns": { "db": "<namespace-db>", "coll": "<namespace-coll>" }, "documentKey": { "_id": { "$oid": "<object-id>" } }, "initialSync": { "phase": "<sync-state>" } "kinesisStream": "<kinesis-name>", "shardId": "<kinesis-shard-id>", "sequenceNumber": "<doc-uuid>", "partitionKey": "<partition-id>", } "window": { "start": <ISODate>, "end": <ISODate>, "partition": "<session-partition>" }, "https": { "url": "<target-url>", "method": "<request-method>", "httpStatusCode": <http-code>, "responseTimeMs": <response-time-ms> }
Sintaxe
A expressão $meta recebe uma única entrada de string que corresponde ao caminho totalmente qualificado em sintaxe de pontos de uma fonte de metadados. A raiz deste caminho deve ser "stream". Você pode query os seguintes caminhos:
Caminho | Tipo | Condicionalidade | Descrição |
|---|---|---|---|
| objeto | Sempre | |
| documento | Sempre | Documento contendo metadados para o estágio |
| string | Sempre | Tipo de conexão usada como origem. |
| Data ISO | Sempre | Data e hora do registro no ponto de ingestão. |
| string | Condicional | Tópico do Kafka do qual o fluxo ingere registros. Aplica-se apenas a uma fonte do Kafka. |
| inteiro | Condicional | Partição do tópico do Kafka de onde o fluxo ingere registros. Aplica-se apenas a uma fonte do Kafka. |
| inteiro | Condicional | Rastreamento de deslocamento da ordem das mensagens e posição na fila dentro de uma partição de origem do Kafka. Aplica-se apenas a uma fonte do Kafka. |
| string|int|long|double|objeto|binData | Condicional | Chave atribuída às mensagens do Kafka para particionamento e distribuição de carga. Aplica-se apenas a uma fonte do Kafka. |
| array | Condicional | Conjunto de pares chave-valor que descrevem os metadados das mensagens do Kafka. Aplica-se apenas a uma fonte do Kafka. |
| string | Condicional | Tipo de operação de banco de dados que o Atlas Stream Processing tentou executar no documento fornecido. Aplica-se apenas a uma fonte de fluxo de alterações do Atlas. |
| documento | Condicional | Documento contendo o namespace do qual o Atlas Stream Processing obtém documentos. Aplica-se apenas a uma fonte de fluxo de alterações do Atlas. |
| string | Condicional | Nome do banco de dados contra o qual o Atlas Stream Processing tenta operações. Aplica-se apenas a uma fonte de fluxo de alterações do Atlas. Esse valor é o mesmo para todos os documentos de uma fonte de fluxo de alterações de coleção ou fluxo de alterações de banco de dados. Varia para uma fonte de fluxo de alterações de cluster. |
| string | Condicional | Nome da coleção contra a qual o Atlas Stream Processing tenta operações. Aplica-se apenas a uma fonte de fluxo de alterações do Atlas. Esse valor é o mesmo para todos os documentos de uma fonte de Coleção Change Stream. Varia para uma fonte de banco de dados Change Stream ou Cluster Change Stream. |
| documento | Condicional | Documento contendo o ID do objeto do documento de origem. Aplica-se apenas a uma fonte de fluxo de alterações do Atlas. |
| string | Condicional | Estado atual da operação de sincronização inicial. Aplica-se apenas a uma fonte de fluxo de alterações do Atlas durante a sincronização inicial. |
| string | Condicional | Nome do Kinesis Data Stream do qual o Atlas Stream Processing obtém documentos. Aplica-se apenas a uma fonte Kinesis do Amazon Web Services. |
| string | Condicional | ID do shard dentro do Kinesis Data Stream do qual o Atlas Stream Processing obtém documentos. Aplica-se apenas a uma fonte Kinesis do Amazon Web Services. |
| string | Condicional | Identificador exclusivo do documento originado do Kinesis Data Stream. Aplica-se apenas a uma fonte Kinesis do Amazon Web Services. |
| string | Condicional | Identificador exclusivo da partição à qual o documento de origem pertence. Aplica-se apenas a uma fonte Kinesis do Amazon Web Services. |
| documento | Condicional | Documento contendo metadados da janela. Aplica-se somente se o documento tiver sido processado em uma janela. |
| Data ISO | Condicional | Tempo de abertura da janela. Aplica-se somente se o documento tiver sido processado em uma janela. |
| Data ISO | Condicional | Tempo de fechamento da janela. Aplica-se somente se o documento tiver sido processado em uma janela. |
| string | Condicional | Partição da janela de sessão à qual o documento pertence. Aplicar somente se o documento tiver sido processado em uma janela de sessão. |
| documento | Condicional | Documento contendo metadados para o estágio $https. Aplica-se somente quando a falha em processamento ocorreu no estágio |
| string | Condicional | URL de destino do estágio |
| string | Condicional | Método de solicitação HTTP usado pelo estágio |
| int | Condicional | Código de status de resposta HTTP da solicitação. Aplica-se apenas quando a falha em processamento ocorreu no estágio |
| int | Condicional | Tempo de resposta da solicitação em milissegundos. Aplica-se apenas quando a falha em processamento ocorreu no estágio |
Comportamento
A expressão $meta do Atlas Stream Processing fornece toda a funcionalidade da expressão de agregação existente do MongoDB $meta. No entanto, você não pode usar a funcionalidade específica para a versão do Atlas Stream Processing do $meta em uma consulta de agregação MongoDB padrão.
Exemplos
O exemplo a seguir enriquece a saída de um fluxo com um array dos tópicos de origem do Kafka dos quais os dados foram ingeridos:
{ $source: { connectionName: "kafka", topic: ["t1", "t2", "t3"] } }, { $emit: { connectionName: "kafka", topic: { $concat: [ { $meta: "stream.source.topic" }, "out" ] } } }
O exemplo a seguir adiciona um campo ao fluxo que relata o relatório da hora de início de cada janela.
{ $source: { connectionName: "kafka", topic: "t1" } }, { $hoppingWindow: . . . }, { $addFields: { start: { $meta: "stream.window.start" } } }