Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
90 commits
Select commit Hold shift + click to select a range
56cfac4
Feat: estruturação do flow de br_rf_cnpj
luizavboas Jul 21, 2026
c1eecf2
Feat: modelos sql e correção de decorator das tasks
luizavboas Jul 21, 2026
886bd0c
Fix: r_run_rf_cnpj
luizavboas Jul 21, 2026
8c08c42
Fix: argumentos main (tabelas->tables)
luizavboas Jul 21, 2026
df1cd65
Fix: utils.py:build_paths
luizavboas Jul 21, 2026
e324eab
Fix: utils.py:build_paths input_path->input_dir
luizavboas Jul 21, 2026
db79a5d
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 21, 2026
687db59
Fix: tasks.py:processamento das tabelas (construção do caminho dos ar…
luizavboas Jul 21, 2026
fca0e9d
git pushMerge branch 'chore/ajustes_cnpj' of github.com:basedosdados/…
luizavboas Jul 21, 2026
6cbf170
Fix: encoding
luizavboas Jul 22, 2026
5b46480
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 22, 2026
4e99cbc
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 22, 2026
2a9aa01
Fix: change chunksoze and add dataset to dbt_project
luizavboas Jul 22, 2026
4da7902
Merge branch 'chore/ajustes_cnpj' of github.com:basedosdados/pipeline…
luizavboas Jul 22, 2026
ccce5d7
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 23, 2026
1a8b233
Test: utils.py alterada para consultar a tabela staging no flow de di…
luizavboas Jul 23, 2026
20c0387
Test: adição do parametro chunk_size ao flow e mudanças no fluxo de d…
luizavboas Jul 23, 2026
820f554
Test: utils.py alterada para consultar a tabela staging no flow de di…
luizavboas Jul 23, 2026
1e51655
Feat: folder date como parâmetro. Adaptações foram feitas em utils.py…
luizavboas Jul 24, 2026
c377969
Fix: ruff check
luizavboas Jul 24, 2026
45c8484
Fix: modificação das colunas de data data_referencia, data_modificaca…
luizavboas Jul 24, 2026
9c6ece4
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 25, 2026
2c52cdb
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 25, 2026
59bd753
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 25, 2026
80e999a
Test: max_parallel de 12 para 8
luizavboas Jul 27, 2026
f76ba4b
Merge branch 'chore/ajustes_cnpj' of github.com:basedosdados/pipeline…
luizavboas Jul 27, 2026
e24ccdf
Test: max_retries, max_parallel, timeout, chunksize as flows parameters
luizavboas Jul 27, 2026
7457ce0
Test: max_retries, max_parallel, timeout, chunksize as flows paramete…
luizavboas Jul 27, 2026
51a7821
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 27, 2026
b88ba3b
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 28, 2026
e607b05
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 28, 2026
81c879c
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 28, 2026
0acf270
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 28, 2026
ee555f9
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Jul 29, 2026
38bb5a0
Fix: header repetido
luizavboas Jul 29, 2026
407d186
Fix: modelos adequados a data_referencia e data_modificacao
luizavboas Aug 3, 2026
fea8aff
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 3, 2026
7e70bed
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 3, 2026
10cdb3f
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 4, 2026
606da16
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 4, 2026
99f378c
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 4, 2026
9ca1b4d
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 4, 2026
8187bd3
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 4, 2026
15e4c5f
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 4, 2026
f9ea693
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 4, 2026
fa71a91
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 5, 2026
bc28e56
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 5, 2026
6d329aa
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 5, 2026
4fb79e5
Test: modelos com dados históricos de br_me_cnpj
luizavboas Aug 5, 2026
23595b1
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 5, 2026
9de529a
Test: flow com o get_table_files como forma de obter links com format…
luizavboas Aug 5, 2026
b275f72
Merge branch 'chore/ajustes_cnpj' of github.com:basedosdados/pipeline…
luizavboas Aug 5, 2026
5097376
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 5, 2026
5a3a5e6
Fix: adequação dos modelos para usar tabelas legado criadas em br_rf_…
luizavboas Aug 5, 2026
8f28fe5
Merge branch 'chore/ajustes_cnpj' of github.com:basedosdados/pipeline…
luizavboas Aug 5, 2026
5e66027
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 5, 2026
3eeddca
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 6, 2026
1876800
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 6, 2026
1c7dcc9
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 6, 2026
cfeeda2
Fix:teste de ingestão de dados 2023-08
luizavboas Aug 6, 2026
10349fb
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 6, 2026
df463df
Merge branch 'chore/ajustes_cnpj' of github.com:basedosdados/pipeline…
luizavboas Aug 6, 2026
f21e95a
Fix: casting de variáveis no flow de dicionário
luizavboas Aug 6, 2026
3754ed8
Fix: billing project na verificação de relacionamentos - flow dicionário
luizavboas Aug 6, 2026
3d4105f
Fix: billing project na verificação de relacionamentos - flow dicionário
luizavboas Aug 6, 2026
00125dc
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 6, 2026
c9d6824
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 7, 2026
9cac3e6
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 7, 2026
a090d6a
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 7, 2026
a1c9115
Feat: adição de __most_recent_date_cnpj__ à macro get_qhere_subquery …
luizavboas Aug 7, 2026
e34573c
git pushMerge branch 'chore/ajustes_cnpj' of github.com:basedosdados/…
luizavboas Aug 7, 2026
248f123
Add: .sql de tabelas legado para migração dos buckets de basedosdados…
luizavboas Aug 7, 2026
2c0c56a
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 7, 2026
78531b7
Fix: pyrefly checks
luizavboas Aug 7, 2026
eb509c1
Merge branch 'chore/ajustes_cnpj' of github.com:basedosdados/pipeline…
luizavboas Aug 7, 2026
cc15a77
Fix: mudança da coluna de dsata em register_table_materialization_task
luizavboas Aug 7, 2026
cfe9aca
Test: remoção dos .sqls de tabelas legado para verificar se o check_m…
luizavboas Aug 7, 2026
4320d76
Test: voltando sqls das tabelas legado
luizavboas Aug 7, 2026
ca9a973
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 7, 2026
5513024
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 8, 2026
49b43e1
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 8, 2026
0bf0309
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 8, 2026
c4fd7d4
Fix: date_format para poll_source_update_task docstring para _run_rf_…
luizavboas Aug 10, 2026
de4a5a6
Fix: date_format de folder_date
luizavboas Aug 10, 2026
3867e86
Test: evitando teste de check metadata
luizavboas Aug 10, 2026
669ac57
Test: evitando teste de check metadata (revert)
luizavboas Aug 10, 2026
d80d5ea
Fix: particionamento simples estava errado e usava data_referencia
luizavboas Aug 10, 2026
04e1750
Fix: particionamento simples estava errado e usava data_referencia
luizavboas Aug 10, 2026
5c1e4ad
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 11, 2026
f89ae45
Merge branch 'main' into chore/ajustes_cnpj
mergify[bot] Aug 11, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions dbt_project.yml
Original file line number Diff line number Diff line change
Expand Up @@ -424,6 +424,9 @@ models:
br_rf_cno:
+materialized: table
+schema: br_rf_cno
br_rf_cnpj:
+materialized: table
+schema: br_rf_cnpj
br_rj_isp_estatisticas_seguranca:
+materialized: table
+schema: br_rj_isp_estatisticas_seguranca
Expand Down
20 changes: 20 additions & 0 deletions macros/custom_get_where_subquery.sql
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,26 @@
{% endif %}
{% endif %}

{# This block looks for __most_recent_date_cnpj__ placeholder #}
{% if "__most_recent_date_cnpj__" in where %}
{% set max_date_query = (
"select max(data_referencia) as max_date from " ~ relation
) %}
{% set max_date_result = run_query(max_date_query) %}
{% if execute and max_date_result.rows[0][0] %}
{% set max_date = max_date_result.rows[0][0] %}
{% set where = where | replace(
"__most_recent_date_cnpj__",
"data_referencia = '" ~ max_date ~ "'",
) %}
{% do log(
"The test will filter by the most recent date: "
~ max_date,
info=True,
) %}
{% endif %}
{% endif %}

{% if "__most_recent_date_cno__" in where %}
{% set max_date_query = (
"select max(data_extracao) as max_date from " ~ relation
Expand Down
62 changes: 62 additions & 0 deletions models/br_rf_cnpj/br_rf_cnpj__dicionario.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
{{ config(alias="dicionario", schema="br_rf_cnpj", materialized="table") }}

with
tmp_dict as (
select
safe_cast(id_tabela as string) id_tabela,
safe_cast(nome_coluna as string) nome_coluna,
safe_cast(chave as string) chave,
safe_cast(cobertura_temporal as string) cobertura_temporal,
regexp_replace(
{{ validate_null_cols("valor") }}, r'(^0+)(?:[^0]+|0{1})', ''
) as valor,
from {{ set_datalake_project("br_rf_cnpj_staging.dicionario") }} as t
)
select
id_tabela,
nome_coluna,
chave,
cobertura_temporal,
case when nome_coluna = "id_pais" then valor else initcap(valor) end as valor
from tmp_dict
where valor is not null
union all
{{
dicionario_not_found(
id_tabela="empresas",
nome_coluna="qualificacao_responsavel",
chave="36",
)
}}
union all
{{
dicionario_not_found(
id_tabela="socios",
nome_coluna="id_pais",
chave=["994", "393"],
)
}}
union all
{{
dicionario_not_found(
id_tabela="estabelecimentos",
nome_coluna="id_pais",
chave=["8", "9", "393"],
)
}}
union all
{{
dicionario_not_found(
id_tabela="estabelecimentos",
nome_coluna="motivo_situacao_cadastral",
chave="32",
)
}}
union all
{{
dicionario_not_found(
id_tabela="estabelecimentos",
nome_coluna="",
chave=["6202100", "4761000"],
)
}}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
55 changes: 55 additions & 0 deletions models/br_rf_cnpj/br_rf_cnpj__empresas.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
{{
config(
schema="br_rf_cnpj",
alias="empresas",
materialized="incremental",
partition_by={
"field": "data_referencia",
"data_type": "date",
"granularity": "month",
},
)
}}

with
cnpj_empresas as (
select
safe.parse_date('%Y-%m', data_referencia) data_referencia,
safe_cast(lpad(cnpj_basico, 8, '0') as string) cnpj_basico,
safe_cast(razao_social as string) razao_social,
safe_cast(natureza_juridica as string) natureza_juridica,
safe_cast(
regexp_replace(qualificacao_responsavel, '^0', '') as string
) qualificacao_responsavel,
safe_cast(capital_social as float64) capital_social,
safe_cast(regexp_replace(porte, '^0', '') as string) porte,
safe_cast(ente_federativo as string) ente_federativo,
safe_cast(data_modificacao as date) data_modificacao
from {{ set_datalake_project("br_rf_cnpj_staging.empresas") }} as t
where
porte != "porte"
{% if is_incremental() %}
and safe.parse_date('%Y-%m', data_referencia)
> (select max(data_referencia) from {{ this }})
{% else %}
-- Dados históricos até 2023-04-30 foram migrados do modelo
-- br_me_cnpj.estabelecimentos
union all
select
safe_cast(data as date) data_referencia,
safe_cast(lpad(cnpj_basico, 8, '0') as string) cnpj_basico,
safe_cast(razao_social as string) razao_social,
safe_cast(natureza_juridica as string) natureza_juridica,
safe_cast(
regexp_replace(qualificacao_responsavel, '^0', '') as string
) qualificacao_responsavel,
safe_cast(capital_social as float64) capital_social,
safe_cast(regexp_replace(porte, '^0', '') as string) porte,
safe_cast(ente_federativo as string) ente_federativo,
safe_cast(null as date) data_modificacao
from {{ set_datalake_project("br_rf_cnpj_staging.empresas_legado") }}
where porte != "porte" and safe_cast(data as date) <= date("2023-04-30")
{% endif %}
)
select *
from cnpj_empresas
Empty file.
148 changes: 148 additions & 0 deletions models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,148 @@
{{
config(
schema="br_rf_cnpj",
alias="estabelecimentos",
materialized="incremental",
partition_by={
"field": "data_referencia",
"data_type": "date",
"granularity": "month",
},
cluster_by=["sigla_uf"],
)
}}
with
cnpj_estabelecimentos as (
select
safe.parse_date('%Y-%m', data_referencia) data_referencia,
safe_cast(lpad(cnpj, 14, "0") as string) cnpj,
safe_cast(lpad(cnpj_basico, 8, '0') as string) cnpj_basico,
safe_cast(lpad(cnpj_ordem, 4, '0') as string) cnpj_ordem,
safe_cast(lpad(cnpj_dv, 2, '0') as string) cnpj_dv,
safe_cast(
identificador_matriz_filial as string
) identificador_matriz_filial,
safe_cast(nome_fantasia as string) nome_fantasia,
safe_cast(cast(situacao_cadastral as int64) as string) situacao_cadastral,
safe_cast(data_situacao_cadastral as date) data_situacao_cadastral,
safe_cast(
regexp_replace(motivo_situacao_cadastral, '^0', '') as string
) motivo_situacao_cadastral,
safe_cast(nome_cidade_exterior as string) nome_cidade_exterior,
safe_cast(cast(id_pais as int64) as string) id_pais,
Comment on lines +26 to +32

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🩺 Stability & Availability | 🟠 Major | ⚡ Quick win

🧩 Analysis chain

🏁 Script executed:

#!/usr/bin/env bash
set -euo pipefail

# Find all nested SAFE_CAST(CAST(...)) expressions in the CNPJ models.
rg -n -U -P 'safe_cast\s*\(\s*cast\s*\(' models/br_rf_cnpj

# Inspect raw CNPJ transformations for handling of blank numeric-code fields.
rg -n -C 4 --glob '*.py' \
  'qualificacao_representante_legal|qualificacao|situacao_cadastral|id_pais' \
  pipelines/crawler/rf_cnpj

Repository: basedosdados/pipelines

Length of output: 12888


🏁 Script executed:

#!/usr/bin/env bash
set -euo pipefail

echo "== model excerpts =="
sed -n '1,100p' models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql
sed -n '1,85p' models/br_rf_cnpj/br_rf_cnpj__socios.sql

echo "== source preparation excerpts around rf_cnpj columns =="
rg -n -C 7 'qualificacao|id_pais|situacao_cadastral|CNPJ|RF_CNPJ|colunas_estabelecimentos|COLUNAS_ESTABELECIMENTOS|COLUNAS_SOCIOS' pipelines/crawler/rf_cnpj -g '*.py'

echo "== deterministic SQL semantics probe if duckdb is available =="
python3 - <<'PY'
from pathlib import Path
p = Path("/tmp/probe_cast.sql")
for val in ("", None, "abc", "1"):
    expr = f"SELECT {repr(val) if val else 'NULL'} AS val, CAST({repr(val) if val else 'NULL'} AS INT64) AS safe_cast_of_cast"
    txt = p.read_text() if p.exists() else ""
    txt += expr + "\n"
    p.write_text(txt)
print("sql generated for values: '', NULL, 'abc', '1'\n" + p.read_text())
PY

Repository: basedosdados/pipelines

Length of output: 31166


🏁 Script executed:

#!/usr/bin/env bash
set -euo pipefail

echo "== rf_cnpj utils relevant source =="
sed -n '430,510p' pipelines/crawler/rf_cnpj/utils.py
sed -n '780,830p' pipelines/crawler/rf_cnpj/utils.py

echo "== deterministic typecast probe if duckdb is available =="
if command -v duckdb >/dev/null 2>&1; then
  cat >/tmp/cast_probe.sql <<'SQL'
select
  'NULL' as value, CAST(NULL AS INT64) as cast_result, CAST(NULL AS STRING) as outer_safe_result
union all select 'empty', CAST('' AS INT64), CAST(CAST('' AS INT64) AS STRING)
union all select 'malformed', CAST('abc' AS INT64), CAST(CAST('abc' AS INT64) AS STRING)
union all select 'valid', CAST('1' AS INT64), CAST(CAST('1' AS INT64) AS STRING);
SQL
  duckdb /tmp/cast_probe.sql
else
  echo "duckdb not available"
fi

Repository: basedosdados/pipelines

Length of output: 5146


Make numeric-code normalization fail-safe.

The inner CAST(... AS INT64) must not run before the outer SAFE_CAST, otherwise blank or malformed numeric-code values fail the model instead of returning NULL.

Replace the inner casts with SAFE_CAST(... AS INT64) at every affected site:

  • models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql: situacao_cadastral and id_pais, including the legacy branch.
  • models/br_rf_cnpj/br_rf_cnpj__socios.sql: qualificacao, id_pais, and qualificacao_representante_legal, including the legacy branch.
📍 Affects 2 files
  • models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql#L26-L32 (this comment)
  • models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql#L73-L81
  • models/br_rf_cnpj/br_rf_cnpj__socios.sql#L21-L28
  • models/br_rf_cnpj/br_rf_cnpj__socios.sql#L47-L54
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql` around lines 26 - 32,
Replace the inner CAST(... AS INT64) with SAFE_CAST(... AS INT64) for
numeric-code normalization in models/br_rf_cnpj/br_rf_cnpj__estabelecimentos.sql
at lines 26-32 and 73-81, covering situacao_cadastral and id_pais in both
branches. Apply the same change in models/br_rf_cnpj/br_rf_cnpj__socios.sql at
lines 21-28 and 47-54 for qualificacao, id_pais, and
qualificacao_representante_legal, preserving the outer casts and NULL behavior
for malformed values.

safe_cast(data_inicio_atividade as date) data_inicio_atividade,
safe_cast(cnae_fiscal_principal as string) cnae_fiscal_principal,
safe_cast(cnae_fiscal_secundaria as string) cnae_fiscal_secundaria,
safe_cast(sigla_uf as string) sigla_uf,
safe_cast(safe_cast(id_municipio_rf as numeric) as string) id_municipio_rf,
safe_cast(tipo_logradouro as string) tipo_logradouro,
safe_cast(logradouro as string) logradouro,
safe_cast(numero as string) numero,
safe_cast(complemento as string) complemento,
safe_cast(bairro as string) bairro,
safe_cast(replace (cep, ".0", "") as string) cep,
safe_cast(ddd_1 as string) ddd_1,
safe_cast(telefone_1 as string) telefone_1,
safe_cast(ddd_2 as string) ddd_2,
safe_cast(telefone_2 as string) telefone_2,
safe_cast(ddd_fax as string) ddd_fax,
safe_cast(fax as string) fax,
safe_cast(lower(email) as string) email,
safe_cast(situacao_especial as string) situacao_especial,
safe_cast(data_modificacao as date) data_modificacao,
safe_cast(data_situacao_especial as date) data_situacao_especial
from {{ set_datalake_project("br_rf_cnpj_staging.estabelecimentos") }}
{% if is_incremental() %}
where
safe.parse_date('%Y-%m', data_referencia)
> (select max(data_referencia) from {{ this }})
-- Dados históricos até 2023-04-30 foram migrados do modelo
-- br_me_cnpj.estabelecimentos
{% else %}
union all
select
safe_cast(data as date) data_referencia,
safe_cast(lpad(cnpj, 14, "0") as string) cnpj,
safe_cast(lpad(cnpj_basico, 8, '0') as string) cnpj_basico,
safe_cast(lpad(cnpj_ordem, 4, '0') as string) cnpj_ordem,
safe_cast(lpad(cnpj_dv, 2, '0') as string) cnpj_dv,
safe_cast(
identificador_matriz_filial as string
) identificador_matriz_filial,
safe_cast(nome_fantasia as string) nome_fantasia,
safe_cast(
cast(situacao_cadastral as int64) as string
) situacao_cadastral,
safe_cast(data_situacao_cadastral as date) data_situacao_cadastral,
safe_cast(
regexp_replace(motivo_situacao_cadastral, '^0', '') as string
) motivo_situacao_cadastral,
safe_cast(nome_cidade_exterior as string) nome_cidade_exterior,
safe_cast(cast(id_pais as int64) as string) id_pais,
safe_cast(data_inicio_atividade as date) data_inicio_atividade,
safe_cast(cnae_fiscal_principal as string) cnae_fiscal_principal,
safe_cast(cnae_fiscal_secundaria as string) cnae_fiscal_secundaria,
safe_cast(sigla_uf as string) sigla_uf,
safe_cast(
safe_cast(id_municipio_rf as numeric) as string
) id_municipio_rf,
safe_cast(tipo_logradouro as string) tipo_logradouro,
safe_cast(logradouro as string) logradouro,
safe_cast(numero as string) numero,
safe_cast(complemento as string) complemento,
safe_cast(bairro as string) bairro,
safe_cast(replace (cep, ".0", "") as string) cep,
safe_cast(ddd_1 as string) ddd_1,
safe_cast(telefone_1 as string) telefone_1,
safe_cast(ddd_2 as string) ddd_2,
safe_cast(telefone_2 as string) telefone_2,
safe_cast(ddd_fax as string) ddd_fax,
safe_cast(fax as string) fax,
safe_cast(lower(email) as string) email,
safe_cast(situacao_especial as string) situacao_especial,
safe_cast(null as date) data_modificacao,
safe_cast(data_situacao_especial as date) data_situacao_especial
from
{{ set_datalake_project("br_rf_cnpj_staging.estabelecimentos_legado") }}
where safe_cast(data as date) <= date("2023-04-30")
{% endif %}
)
select
a.data_referencia,
a.cnpj,
a.cnpj_basico,
a.cnpj_ordem,
a.cnpj_dv,
a.identificador_matriz_filial,
a.nome_fantasia,
a.situacao_cadastral,
a.data_situacao_cadastral,
a.motivo_situacao_cadastral,
a.nome_cidade_exterior,
a.id_pais,
a.data_inicio_atividade,
a.cnae_fiscal_principal,
a.cnae_fiscal_secundaria,
a.sigla_uf,
safe_cast(b.id_municipio as string) id_municipio,
a.id_municipio_rf,
a.tipo_logradouro,
a.logradouro,
a.numero,
a.complemento,
a.bairro,
a.cep,
a.ddd_1,
a.telefone_1,
a.ddd_2,
a.telefone_2,
a.ddd_fax,
a.fax,
a.email,
a.situacao_especial,
a.data_modificacao,
a.data_situacao_especial
from cnpj_estabelecimentos a
left join
basedosdados.br_bd_diretorios_brasil.municipio b
on safe_cast(safe_cast(a.id_municipio_rf as numeric) as string) = b.id_municipio_rf
Empty file.
20 changes: 20 additions & 0 deletions models/br_rf_cnpj/br_rf_cnpj__simples.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
{{
config(
schema="br_rf_cnpj",
alias="simples",
materialized="table",
)
}}

select
-- safe.parse_date('%Y-%m', data_referencia) data_referencia,
lpad(safe_cast(cnpj_basico as string), 8, '0') cnpj_basico,
safe_cast(opcao_simples as int64) opcao_simples,
safe_cast(data_opcao_simples as date) data_opcao_simples,
safe_cast(data_exclusao_simples as date) data_exclusao_simples,
safe_cast(opcao_mei as int64) opcao_mei,
safe_cast(data_opcao_mei as date) data_opcao_mei,
safe_cast(data_exclusao_mei as date) data_exclusao_mei,
-- safe_cast(data_modificacao as date) data_modificacao
from {{ set_datalake_project("br_rf_cnpj_staging.simples") }} as t
where safe_cast(opcao_mei as string) != "opcao_mei"
Empty file.
64 changes: 64 additions & 0 deletions models/br_rf_cnpj/br_rf_cnpj__socios.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
{{
config(
schema="br_rf_cnpj",
alias="socios",
materialized="incremental",
partition_by={
"field": "data_referencia",
"data_type": "date",
"granularity": "month",
},
)
}}
with
cnpj_socios as (
select
safe.parse_date('%Y-%m', data_referencia) data_referencia,
lpad(safe_cast(cnpj_basico as string), 8, '0') cnpj_basico,
safe_cast(tipo as string) tipo,
safe_cast(nome as string) nome,
safe_cast(documento as string) documento,
safe_cast(cast(qualificacao as int64) as string) qualificacao,
safe_cast(data_entrada_sociedade as date) data_entrada_sociedade,
safe_cast(cast(id_pais as int64) as string) id_pais,
safe_cast(cpf_representante_legal as string) cpf_representante_legal,
safe_cast(nome_representante_legal as string) nome_representante_legal,
safe_cast(
cast(qualificacao_representante_legal as int64) as string
) qualificacao_representante_legal,
safe_cast(faixa_etaria as string) faixa_etaria,
safe_cast(data_modificacao as date) data_modificacao
from {{ set_datalake_project("br_rf_cnpj_staging.socios") }} as t
where
safe_cast(qualificacao as string) != "qualificacao"
{% if is_incremental() %}
and safe.parse_date('%Y-%m', data_referencia)
> (select max(data_referencia) from {{ this }})
{% else %}
-- Dados históricos até 2023-04-30 foram migrados do modelo
-- br_me_cnpj.socios
union all
select
safe_cast(data as date) data_referencia,
lpad(safe_cast(cnpj_basico as string), 8, '0') cnpj_basico,
safe_cast(tipo as string) tipo,
safe_cast(nome as string) nome,
safe_cast(documento as string) documento,
safe_cast(cast(qualificacao as int64) as string) qualificacao,
safe_cast(data_entrada_sociedade as date) data_entrada_sociedade,
safe_cast(cast(id_pais as int64) as string) id_pais,
safe_cast(cpf_representante_legal as string) cpf_representante_legal,
safe_cast(nome_representante_legal as string) nome_representante_legal,
safe_cast(
cast(qualificacao_representante_legal as int64) as string
) qualificacao_representante_legal,
safe_cast(faixa_etaria as string) faixa_etaria,
safe_cast(null as date) data_modificacao
from {{ set_datalake_project("br_rf_cnpj_staging.socios_legado") }}
where
safe_cast(qualificacao as string) != "qualificacao"
and safe_cast(data as date) <= date("2023-04-30")
{% endif %}
)
select *
from cnpj_socios
Loading
Loading