Definição
O estágio $lookup executa uma junção externa esquerda do fluxo de mensagens de sua $source para uma collection do Atlas em seu Registro de Conexão.
Dependendo do seu caso de uso, um estágio de pipeline $lookup usa uma das três sintaxes a seguir:
Para saber mais, consulte Sintaxe $lookup.
Aviso
Usar $lookup para enriquecimento de um fluxo pode reduzir a velocidade de processamento do fluxo.
O seguinte formulário protótipo ilustra todos os campos disponíveis:
{ "$lookup": { "from": { "connectionName": "<registered-atlas-connection>", "db": "<registered-database-name>", "coll": "<atlas-collection-name>", "readPreference": "<read-preference>", "readPreferenceTags": [ {"<key>": "<value>"}, ... ] }, "localField": "<field-in-source-messages>", "foreignField": "<field-in-from-collection>", "let": { <var_1>: <expression>, <var_2>: <expression>, …, <var_n>: <expression> }, "pipeline": [ <pipeline to run> ], "as": "<output-array-field>", "parallelism": <integer>, "partitionBy": <expression> } }
Sintaxe
O estágio $lookup recebe um documento com os seguintes campos:
Campo | Tipo | necessidade | Descrição |
|---|---|---|---|
from | documento | Condicional | Documento que especifica uma collection em um banco de dados Atlas para unir às mensagens do seu Se você especificar este campo, você deverá especificar valores para todos os campos neste documento. Este campo não é exigido se você especificar um campo |
from.connectionName | string | Condicional | Nome da conexão no registro de conexões. Este campo não é exigido se você especificar um campo |
from.db | string | Condicional | Nome do banco de dados do Atlas que contém a coleção que você deseja unir. Este campo não é exigido se você especificar um campo |
from.coll | string | Condicional | Nome da coleção da qual você deseja participar. Este campo não é exigido se você especificar um campo |
from.readPreference | string | Opcional | Read preference para Padrão é |
from.readPreferenceTags | array | Opcional | Leia as marcações de preferência para |
localField | string | Condicional | Campo a partir de suas mensagens Este campo faz parte das seguintes sintaxes: |
foreignField | string | Condicional | Campo de documentos na coleção Este campo faz parte das seguintes sintaxes: |
let | documento | Condicional | |
gasoduto | documento | Condicional | Especifica o Este campo faz parte das seguintes sintaxes: |
como | string | Obrigatório | Nome do novo campo de array a ser adicionado aos documentos de entrada. Este novo campo de array contém os documentos correspondentes da coleção |
| inteiro | Opcional | Número máximo de solicitações paralelas feitas para o destino Deve ser um número inteiro entre Cada processador de fluxo tem um valor máximo de paralelismo cumulativo determinado por seu nível. O paralelismo cumulativo de um processador de fluxo é calculado da seguinte forma:
Em que Por exemplo, se o seu estágio Se um processador de fluxo exceder o paralelismo cumulativo máximo para seu nível, o Atlas Stream Processing lançará um erro e o avisará sobre o nível mínimo de processador necessário para o nível pretendido de paralelismo. Você deve dimensionar o processador para um nível superior ou reduzir os valores de paralelismo de seus estágios para resolver o erro. Para saber mais, consulte Stream Processing. |
| expressão | Opcional | Expressão usada para dividir o fluxo de entrada em threads paralelos. O Atlas Stream Processing aloca cada documento de entrada com o mesmo resultado de expressão Se este campo não for especificado, os documentos de entrada serão enviados para threads paralelos por round-robin. |
Comportamento
A versão Atlas Stream Processing de $lookup executa uma junção externa esquerda de mensagens de seu $source e os documentos em uma collection Atlas especificada. Essa versão se comporta de forma semelhante ao estágio disponível em um banco de dados MongoDB padrão.$lookup No entanto, esta versão exige que você especifique uma coleção do Atlas a partir do seu Registro de Conexão como o valor para o from campo.
O pipeline pode conter um estágio$lookupaninhado. Se você incluir um estágio$lookupaninhado em seu pipeline, deverá usar a sintaxe from padrão para especificar uma coleção na mesma conexão do Atlas remoto que o estágio$lookupexterno.
Exemplo
$lookup : { from: {connectionName: "dbsrv1", db: "db1", coll: "coll1"}, …, pipeline: [ …, { $lookup: { from: "coll2", …, } }, …, ] }
Se o seu pipeline tiver$lookupe$mergena mesma coleção, os resultados do Atlas Stream Processing poderão variar se você tentar manter uma exibição incremental. O Atlas Stream Processing processa várias mensagens de origem simultaneamente e, em seguida, as mescla todas. Se várias mensagens tiverem o mesmo ID, que$lookupe$mergeusam, o Atlas Stream Processing poderá retornar resultados que ainda não foram materializados.
Exemplo
Considere o seguinte fluxo de entrada:
{ _id: 1, count: 2 } { _id: 1, count: 3 }
Suponha que sua query contenha o seguinte dentro do pipeline:
{ ..., pipeline: [ { $lookup on _id == foreignDoc._id from collection A } { $project: { _id: 1, count: $count + $foreignDoc.count } } { $merge: { into collection A } } ] }
Se estiver tentando manter uma visualização gradual, o resultado poderá ser semelhante ao seguinte:
{ _id: 1, count: 5 }
No entanto, o Atlas Stream Processing pode retornar uma contagem de 5 ou 3, dependendo se o Atlas Stream Processing processou os documentos.
Para mais informações,$lookup consulte.
Exemplos
Uma fonte de dados de streaming gera relatórios meteorológicos detalhados a partir de vários locais, em conformidade com o esquema do Conjunto de dados meteorológicos de amostra. Uma coleção chamada humidity_descriptions contém documentos no formato:
Onde o campo relative_humidity descreve a umidade relativa em temperatura ambiente (20 graus Celsius), e condition lista descritores verbais adequados para esse nível de umidade. Você pode usar a fase $lookup para enriquecer os relatórios meteorológicos em tempo real com descritores sugeridos para meteorologistas usarem nas transmissões meteorológicas.
A seguinte agregação tem quatro fases:
O estágio estabelece uma
$sourceconexão com o intermediário do Apache Kafka, coletando esses relatórios em ummy_weatherdatatópico chamado, expondo cada registro à medida que ele é ingerido para os estágios de agregação subsequentes. Esse estágio também substitui o nome do campo de carimbo de data/hora que ele projeta, definindo-oingestionTimecomo.O estágio
$lookupjunta os registros do banco de dadoshumidity_descriptionsaos relatórios meteorológicos no campodewPoint.O estágio exclui documentos que
$matchtêmhumidity_infoum campo vazio e passa documentos com umhumidity_infocampo preenchido para o próximo estágio.O estágio escreve o resultado em uma collection do Atlas
$mergechamadaenriched_streamnosample_weatherstreambanco de dados. Se não existir tal banco de dados ou coleção, o Atlas os criará.
{ '$source': { connectionName: 'sample_weatherdata', topic: 'my_weatherdata', tsFieldName: 'ingestionTime' } }, { '$lookup': { from: { connectionName: 'weatherStream', db: 'humidity', coll: 'humidity_descriptions' }, 'localField':'dewPoint.value', 'foreignField':'dewPoint', 'as': 'humidity_info' } } { '$match': { 'humidity_info': { '$ne': [] } } } { '$merge': { into: { connectionName: 'weatherStream', db: 'sample_weatherstream', coll: 'enriched_stream' } } }
Para visualizar os documentos na coleção sample_weatherstream.enriched_stream resultante, conecte-se ao cluster Atlas e execute o seguinte comando:
db.getSiblingDB("sample_weatherstream").enriched_stream.find()
{ st: 'x+55100+006100', position: { type: 'Point', coordinates: [ 92.7, -53.6 ] }, elevation: 9999, callLetters: 'UECN', qualityControlProcess: 'V020', dataSource: '4', type: 'FM-13', airTemperature: { value: -11, quality: '9' }, dewPoint: { value: 12.5, quality: '1' }, pressure: { value: 1032.7, quality: '9' }, wind: { direction: { angle: 300, quality: '9' }, type: '9', speed: { rate: 23.6, quality: '2' } }, visibility: { distance: { value: 14000, quality: '1' }, variability: { value: 'N', quality: '1' } }, skyCondition: { ceilingHeight: { value: 390, quality: '9', determination: 'C' }, cavok: 'N' }, sections: [ 'SA1', 'AA1', 'OA1', 'AY1', 'AG1' ], precipitationEstimatedObservation: { discrepancy: '4', estimatedWaterDepth: 21 }, atmosphericPressureChange: { tendency: { code: '1', quality: '1' }, quantity3Hours: { value: 5.5, quality: '1' }, quantity24Hours: { value: 99.9, quality: '9' } }, seaSurfaceTemperature: { value: 1.3, quality: '9' }, waveMeasurement: { method: 'M', waves: { period: 4, height: 2.5, quality: '9' }, seaState: { code: '00', quality: '9' } }, pastWeatherObservationManual: { atmosphericCondition: { value: '4', quality: '1' }, period: { value: 6, quality: '1' } }, skyConditionObservation: { totalCoverage: { value: '07', opaque: '99', quality: '1' }, lowestCloudCoverage: { value: '06', quality: '1' }, lowCloudGenus: { value: '07', quality: '9' }, lowestCloudBaseHeight: { value: 2250, quality: '9' }, midCloudGenus: { value: '07', quality: '9' }, highCloudGenus: { value: '00', quality: '1' } }, presentWeatherObservationManual: { condition: '75', quality: '1' }, atmosphericPressureObservation: { altimeterSetting: { value: 9999.9, quality: '9' }, stationPressure: { value: 1032.6, quality: '1' } }, skyCoverLayer: { coverage: { value: '09', quality: '1' }, baseHeight: { value: 240, quality: '9' }, cloudType: { value: '99', quality: '9' } }, liquidPrecipitation: { period: 6, depth: 3670, condition: '9', quality: '9' }, extremeAirTemperature: { period: 99.9, code: 'N', value: -30.9, quantity: '9' }, ingestionTime: ISODate('2024-09-19T20:04:34.346Z'), humidity_info: [ { _id: ObjectId('66ec805ad3cfbba767ebf7a5'), dewPoint: 12.5, relativeHumidity: 62, condition: 'humid, muggy' } ], }
Observação
O exemplo anterior é representativo. Os dados de streaming não são estáticos e cada usuário vê documentos distintos.