added database interaction class
This commit is contained in:
@@ -0,0 +1,75 @@
|
|||||||
|
{
|
||||||
|
"cells": [
|
||||||
|
{
|
||||||
|
"cell_type": "code",
|
||||||
|
"execution_count": 2,
|
||||||
|
"metadata": {},
|
||||||
|
"outputs": [
|
||||||
|
{
|
||||||
|
"data": {
|
||||||
|
"text/plain": [
|
||||||
|
"{'username': '@otwndk',\n",
|
||||||
|
" 'time': datetime.datetime(2022, 12, 16, 21, 19, 9, tzinfo=tzutc()),\n",
|
||||||
|
" 'content': 'Sponsoreret: Verdensmestre, europamestre og store danske medaljetagere rammer Odense til et prestigefyldt stævne\\n#odense #fyn #cykling'}"
|
||||||
|
]
|
||||||
|
},
|
||||||
|
"execution_count": 2,
|
||||||
|
"metadata": {},
|
||||||
|
"output_type": "execute_result"
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"source": [
|
||||||
|
"from database import Data_table\n",
|
||||||
|
"from scraper import scrape\n",
|
||||||
|
"\n",
|
||||||
|
"topic = 'Odense'\n",
|
||||||
|
"username_list, timestamp_list, content_list = scrape(topic, limit=3, headless=False)\n",
|
||||||
|
"\n",
|
||||||
|
"val_dict = {\n",
|
||||||
|
" 'username': username_list[0],\n",
|
||||||
|
" 'time': timestamp_list[0],\n",
|
||||||
|
" 'content': content_list[0]\n",
|
||||||
|
"}\n",
|
||||||
|
"\n",
|
||||||
|
"print(val_dict)\n",
|
||||||
|
"\n",
|
||||||
|
"with Data_table() as dt:\n",
|
||||||
|
" dt.insert(val_dict)"
|
||||||
|
]
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"cell_type": "code",
|
||||||
|
"execution_count": null,
|
||||||
|
"metadata": {},
|
||||||
|
"outputs": [],
|
||||||
|
"source": []
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"metadata": {
|
||||||
|
"kernelspec": {
|
||||||
|
"display_name": "venv",
|
||||||
|
"language": "python",
|
||||||
|
"name": "python3"
|
||||||
|
},
|
||||||
|
"language_info": {
|
||||||
|
"codemirror_mode": {
|
||||||
|
"name": "ipython",
|
||||||
|
"version": 3
|
||||||
|
},
|
||||||
|
"file_extension": ".py",
|
||||||
|
"mimetype": "text/x-python",
|
||||||
|
"name": "python",
|
||||||
|
"nbconvert_exporter": "python",
|
||||||
|
"pygments_lexer": "ipython3",
|
||||||
|
"version": "3.10.6"
|
||||||
|
},
|
||||||
|
"orig_nbformat": 4,
|
||||||
|
"vscode": {
|
||||||
|
"interpreter": {
|
||||||
|
"hash": "1d532f3642617da00cbdbdc5ef98bd8d4ca226c95ac593ade79d8cbc880a2afe"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"nbformat": 4,
|
||||||
|
"nbformat_minor": 2
|
||||||
|
}
|
||||||
+150
@@ -0,0 +1,150 @@
|
|||||||
|
#!/Users/brian/Projects/twitter_scraper/venv/bin/python
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import traceback
|
||||||
|
import urllib.parse
|
||||||
|
|
||||||
|
import pandas as pd
|
||||||
|
import psycopg2
|
||||||
|
from dotenv import load_dotenv
|
||||||
|
from sqlalchemy import create_engine
|
||||||
|
|
||||||
|
|
||||||
|
class Database(object):
|
||||||
|
def __init__(self):
|
||||||
|
# generate uri
|
||||||
|
engine = os.getenv('DATABASE_ENGINE')
|
||||||
|
username = os.getenv('DATABASE_USERNAME')
|
||||||
|
password = os.getenv('DATABASE_PASSWORD')
|
||||||
|
password = urllib.parse.quote_plus(password) # escape special characters
|
||||||
|
host = os.getenv('DATABASE_HOST')
|
||||||
|
port = os.getenv('DATABASE_PORT')
|
||||||
|
database = os.getenv('DATABASE_DATABASE')
|
||||||
|
self._uri = f'{engine}://{username}:{password}@{host}:{port}/{database}'
|
||||||
|
# prepare connection to database
|
||||||
|
self._engine = create_engine(self._uri)
|
||||||
|
self.con = None
|
||||||
|
|
||||||
|
def __enter__(self):
|
||||||
|
self.connect()
|
||||||
|
return self
|
||||||
|
|
||||||
|
def __exit__(self, exc_type, exc_value, tb):
|
||||||
|
if not exc_type is None:
|
||||||
|
traceback.print_exception(exc_type, exc_value, tb)
|
||||||
|
self.disconnect()
|
||||||
|
|
||||||
|
def connect(self):
|
||||||
|
try:
|
||||||
|
con = self._engine.connect()
|
||||||
|
except psycopg2.OperationalError as e:
|
||||||
|
logging.warning(e)
|
||||||
|
else:
|
||||||
|
self.con = con
|
||||||
|
|
||||||
|
def disconnect(self):
|
||||||
|
self.con.close()
|
||||||
|
|
||||||
|
def execute(self, query, var=None):
|
||||||
|
assert self.con is not None, 'connection to database could not be established'
|
||||||
|
try:
|
||||||
|
if var is None:
|
||||||
|
self.con.execute(query)
|
||||||
|
else:
|
||||||
|
self.con.execute(query, var)
|
||||||
|
except Exception as e:
|
||||||
|
print(e)
|
||||||
|
raise
|
||||||
|
|
||||||
|
def get(self, query):
|
||||||
|
df = pd.read_sql_query(query, self.con)
|
||||||
|
return df
|
||||||
|
|
||||||
|
class Schema(Database):
|
||||||
|
def __init__(self, name):
|
||||||
|
super(Schema, self).__init__()
|
||||||
|
self._name = name
|
||||||
|
|
||||||
|
def __enter__(self):
|
||||||
|
self.connect()
|
||||||
|
if not self.exists():
|
||||||
|
self.create()
|
||||||
|
return self
|
||||||
|
|
||||||
|
def exists(self):
|
||||||
|
query = (
|
||||||
|
'SELECT schema_name '
|
||||||
|
'FROM information_schema.schemata '
|
||||||
|
'ORDER BY schema_name '
|
||||||
|
';'
|
||||||
|
)
|
||||||
|
df = self.get(query)
|
||||||
|
schema_list = df['schema_name'].tolist()
|
||||||
|
return self._name in schema_list
|
||||||
|
|
||||||
|
def create(self):
|
||||||
|
query = f'CREATE SCHEMA {self._name};'
|
||||||
|
self.execute(query)
|
||||||
|
|
||||||
|
def purge(self):
|
||||||
|
query = f'DROP SCHEMA IF EXISTS {self._name} CASCADE;'
|
||||||
|
self.execute(query)
|
||||||
|
self.create()
|
||||||
|
|
||||||
|
class Table(Database):
|
||||||
|
def __init__(self, name):
|
||||||
|
super(Table, self).__init__()
|
||||||
|
self._name = name
|
||||||
|
|
||||||
|
def __enter__(self):
|
||||||
|
self.connect()
|
||||||
|
if not self.exists():
|
||||||
|
self.create()
|
||||||
|
return self
|
||||||
|
|
||||||
|
def exists(self):
|
||||||
|
query = (
|
||||||
|
'SELECT table_name '
|
||||||
|
'FROM information_schema.tables '
|
||||||
|
"WHERE table_schema='public' "
|
||||||
|
'ORDER BY table_name '
|
||||||
|
)
|
||||||
|
df = self.get(query)
|
||||||
|
table_list = df['table_name'].tolist()
|
||||||
|
return self._name in table_list
|
||||||
|
|
||||||
|
def create(self):
|
||||||
|
# placeholder for childrens' methods
|
||||||
|
pass
|
||||||
|
|
||||||
|
def purge(self):
|
||||||
|
query = f'DROP TABLE IF EXISTS {self._name} CASCADE;'
|
||||||
|
self.execute(query)
|
||||||
|
self.create()
|
||||||
|
|
||||||
|
class Data_table(Table):
|
||||||
|
def __init__(self):
|
||||||
|
super(Data_table, self).__init__('data_table')
|
||||||
|
|
||||||
|
def create(self):
|
||||||
|
query = (
|
||||||
|
f'CREATE TABLE {self._name} '
|
||||||
|
'('
|
||||||
|
'id SERIAL PRIMARY KEY, '
|
||||||
|
'username VARCHAR(15) NOT NULL, '
|
||||||
|
'time TIMESTAMPTZ NOT NULL, '
|
||||||
|
'content TEXT NOT NULL, '
|
||||||
|
'UNIQUE (username, time, content) '
|
||||||
|
')'
|
||||||
|
';'
|
||||||
|
)
|
||||||
|
self.execute(query)
|
||||||
|
|
||||||
|
def insert(self, val_dict):
|
||||||
|
query = (
|
||||||
|
f'INSERT INTO {self._name} (username, time, content) '
|
||||||
|
'VALUES (%(username)s, %(time)s, %(content)s) '
|
||||||
|
'ON CONFLICT (username, time, content) DO NOTHING '
|
||||||
|
';'
|
||||||
|
)
|
||||||
|
self.execute(query, val_dict)
|
||||||
Reference in New Issue
Block a user