airbytehq / airbytehq/airbyte

New Integration Request: URL Tables

Đang mở
#3,052 2 bình luận 0 reaction 0 người được giao Xem trên GitHub
area/connectors autoteam Icebox new-connector team/extensibility team/marketplace
Ngôn ngữ chính
Python
Star
22.1k
Fork
5.3k
Chỉ số merge pull request
Chỉ số pull request đang chờ

Mô tả

## Tell us about the new integration you’d like to have
* Which source and which destination?
* Do you need a specific version of the underlying data source e.g: you specifically need support for an older version of the API or DB?

OK this is a slightly strange one, however I often seem to need to snapshot table data from public URLs and stream the content as JSON into BigQuery.

## Describe the context around this new integration
* Which team in your company wants this integration, what for? This helps us understand the use case.
* How often do you want to run syncs?
* If this is an API source connector, which entities/endpoints do you need supported?

Typically I need to take daily snapshots for a historical record of data which is not available via any API.

## Describe the alternative you are considering or using
What are you considering doing if you don’t have this integration through Airbyte?

I have tested this code in a notebook and it works fine, and I could spin it up into a scheduled Cloud Function (triggered by PubSub with the config in the payload), however I am keen on trying to use Airbyte as my single point of data ingestion so I thought this might be an interesting lightweight connector which could enable Airbyte to be used as a (limited) very simple web table scraper:

```
import ssl
import os
import json
import pandas as pd
from datetime import datetime
from google.cloud import bigquery

# specific to this URL due to SSL verification issues
ssl._create_default_https_context = ssl._create_unverified_context

destination_ref = 'beepbeeptechnology.urltable.granada_covid_stats'
source_url = "https://www.juntadeandalucia.es/institutodeestadisticaycartografia/salud/static/resultadosProvincialesCovid_18.html5?prov=18&CodOper=b3_2314&codConsulta=38667"

response = pd.read_html(source_url,encoding='UTF-8')

for response_table in response:
response_table_json = json.dumps(response_table.to_dict(orient='records'), ensure_ascii=False)
stream_json_into_bq_with_id(source_url, response_table_json, destination_ref)
```

Note that this depends on the following function:

```
def stream_json_into_bq_with_id(id_string: str, json_string: str, destination_ref: str) -> dict:
"""
Streams a JSON string as text into a single row in a three column (all string) BigQuery table for subsequent decoding, with the following schema:

Args:
id_string (string): row identifier
json_string (string): data payload
destination_ref (string): BigQuery table destination reference (project_id.dataset_id.table_name)

Returns:
A dict containing
status (string): "success" for successful stream, "error" for a failed stream or "fail" for other failure.
message (string): Exception or error message if applicable
"""
try:
current_time = datetime.now()
current_timestamp = current_time.strftime("%Y-%m-%d %H:%M:%S.%f")

project_id = os.environ.get('GCP_PROJECT')
BQ = bigquery.Client(project=project_id)
table = BQ.get_table(destination_ref)

rows_to_insert = [(id_string, json_string, current_timestamp)]

errors = BQ.insert_rows(
table, rows_to_insert, row_ids=[None] * len(rows_to_insert)
)

stream_log: dict
if len(errors) == 0:
stream_log = {"status": "success", "message": "no errors"}
else:
error_string = json.dumps(errors)
stream_log = {"status": "error", "message": error_string}

except Exception as e:
stream_log = {"status": "fail", "message": e}

return stream_log
```

┆Issue is synchronized with this [Asana task](https://app.asana.com/0/1200367912513076/1200368162572400) by [Unito](https://www.unito.io)

Hướng dẫn đóng góp

Mở hướng dẫn đóng góp

Đánh giá

Issue này chưa được đánh giá.

Nhận issue mới trong hộp thư của bạn

Bản tóm tắt ngắn những issue GitHub phù hợp với người mới.