Definición
La expresión $meta devuelve un objeto que contiene todos los metadatos de transmisión para un documento. Puede exponer estos datos para todo el flujo, o bien para una de las siguientes etapas de agregación de Atlas Stream Processing:
Una expresión $meta tiene la siguiente forma de prototipo:
{ "$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> }
Sintaxis
La expresión $meta toma una única string como entrada, que corresponde a la ruta totalmente calificada en sintaxis de puntos de una fuente de metadatos. La raíz de esta ruta debe ser "stream". Puedes query los siguientes caminos:
ruta | Tipo | Condicionalidad | Descripción |
|---|---|---|---|
| Objeto | Siempre | |
| Documento | Siempre | Documento que contiene metadatos para la |
| string | Siempre | Tipo de conexión utilizada como fuente. |
| ISODate | Siempre | Fecha y hora del registro en el momento de la ingestión. |
| string | Condicional | Tema de Kafka del que el flujo ingiere registros. Se aplica solo a una fuente de Kafka. |
| entero | Condicional | Partición del tema de Kafka del que el flujo ingiere registros. Se aplica solo a una fuente de Kafka. |
| entero | Condicional | Seguimiento del orden de los mensajes y la posición en la cola dentro de una partición de origen de Kafka. Se aplica solo a una fuente de Kafka. |
| string|int|long|double|objeto|binData | Condicional | Clave asignada a los mensajes de Kafka para el particionamiento y la distribución de la carga. Se aplica solo a una fuente de Kafka. |
| arreglo | Condicional | Conjunto de pares clave-valor que describen los metadatos del mensaje de Kafka. Se aplica solo a una fuente de Kafka. |
| string | Condicional | Tipo de operación de base de datos que Atlas Stream Processing intentó realizar en el documento dado. Se aplica solo a una fuente de flujo de cambios de Atlas. |
| Documento | Condicional | Documento que contiene el namespace del que Atlas Stream Processing obtiene documentos. Se aplica solo a una fuente de flujo de cambios de Atlas. |
| string | Condicional | Nombre de la base de datos en la que Atlas Stream Processing intenta realizar operaciones. Se aplica solo a una fuente de flujo de cambios de Atlas. Este valor es el mismo para todos los documentos de una fuente de flujo de cambios de colección o de flujo de cambios de base de datos. Varía para una fuente de flujo de cambios de clúster. |
| string | Condicional | Nombre de la colección contra la que Atlas Stream Processing intenta realizar operaciones. Se aplica solo a una fuente de flujo de cambios de Atlas. Este valor es el mismo para todos los documentos de una fuente de Collection Change Stream. Varía para una fuente de Change Stream de base de datos o Change Stream de clúster. |
| Documento | Condicional | Documento que contiene el ID de objeto del documento de origen. Se aplica solo a una fuente de flujo de cambios de Atlas. |
| string | Condicional | Estado actual de la operación de sincronización inicial. Se aplica solo a una fuente de flujo de cambios de Atlas durante la sincronización inicial. |
| string | Condicional | Nombre del Kinesis Data Stream del que Atlas Stream Processing obtiene documentos. Se aplica solo a una fuente de Kinesis de AWS. |
| string | Condicional | ID de la partición dentro de Kinesis Data Stream desde la que Atlas Stream Processing obtiene los documentos. Solo se aplica a una fuente de Kinesis de Amazon Web Services. |
| string | Condicional | Identificador único del documento obtenido de Kinesis Data Stream. Solo se aplica a una fuente de Kinesis de Amazon Web Services. |
| string | Condicional | Identificador único de la partición a la que pertenece el documento de origen. Solo se aplica a una fuente de Kinesis de Amazon Web Services. |
| Documento | Condicional | Documento que contiene metadatos de la ventana. Solo se aplica si el documento se procesó en una ventana. |
| ISODate | Condicional | Tiempo de apertura de la ventana. Solo se aplica si el documento se procesó en una ventana. |
| ISODate | Condicional | Hora de cierre de la ventana. Solo se aplica si el documento se procesó en una ventana. |
| string | Condicional | Partición de la ventana de sesión a la que pertenece el documento. Solo se aplica si el documento se procesó en una ventana de sesión. |
| Documento | Condicional | Documento que contiene metadatos para la etapa $https. Solo se aplica cuando el error de procesamiento se produjo en la etapa |
| string | Condicional | URL de destino de la etapa |
| string | Condicional | Método de solicitud HTTP utilizado por la etapa |
| Int | Condicional | Código de estado de la respuesta HTTP de la solicitud. Se aplica únicamente cuando se produjo el error de procesamiento en la fase |
| Int | Condicional | Tiempo de respuesta de la solicitud en milisegundos. Se aplica únicamente cuando se produjo el error de procesamiento en la fase |
Comportamiento
La expresión $meta de Stream Processing de Atlas proporciona toda la funcionalidad de la expresión de agregación de MongoDB existente $meta. Sin embargo, no se puede utilizar la funcionalidad específica de la versión de Atlas Stream Processing de $meta en una query estándar de agregación de MongoDB.
Ejemplos
El siguiente ejemplo enriquece la salida de un flujo con un arreglo de los temas fuente de Kafka desde los cuales se ingresaron los datos:
{ $source: { connectionName: "kafka", topic: ["t1", "t2", "t3"] } }, { $emit: { connectionName: "kafka", topic: { $concat: [ { $meta: "stream.source.topic" }, "out" ] } } }
El siguiente ejemplo añade un campo al stream que indica la hora de inicio de cada ventana.
{ $source: { connectionName: "kafka", topic: "t1" } }, { $hoppingWindow: . . . }, { $addFields: { start: { $meta: "stream.window.start" } } }