update
Browse files
app.py
CHANGED
|
@@ -6,39 +6,9 @@ import numpy as np
|
|
| 6 |
import pandas as pd
|
| 7 |
from fastapi import FastAPI
|
| 8 |
from gamma.utils import association
|
| 9 |
-
from kafka import KafkaProducer
|
| 10 |
from pydantic import BaseModel
|
| 11 |
from pyproj import Proj
|
| 12 |
|
| 13 |
-
# Kafak producer
|
| 14 |
-
use_kafka = False
|
| 15 |
-
|
| 16 |
-
try:
|
| 17 |
-
print("Connecting to k8s kafka")
|
| 18 |
-
BROKER_URL = "quakeflow-kafka-headless:9092"
|
| 19 |
-
# BROKER_URL = "34.83.137.139:9094"
|
| 20 |
-
producer = KafkaProducer(
|
| 21 |
-
bootstrap_servers=[BROKER_URL],
|
| 22 |
-
key_serializer=lambda x: dumps(x).encode("utf-8"),
|
| 23 |
-
value_serializer=lambda x: dumps(x).encode("utf-8"),
|
| 24 |
-
)
|
| 25 |
-
use_kafka = True
|
| 26 |
-
print("k8s kafka connection success!")
|
| 27 |
-
except BaseException:
|
| 28 |
-
print("k8s Kafka connection error")
|
| 29 |
-
try:
|
| 30 |
-
print("Connecting to local kafka")
|
| 31 |
-
producer = KafkaProducer(
|
| 32 |
-
bootstrap_servers=["localhost:9092"],
|
| 33 |
-
key_serializer=lambda x: dumps(x).encode("utf-8"),
|
| 34 |
-
value_serializer=lambda x: dumps(x).encode("utf-8"),
|
| 35 |
-
)
|
| 36 |
-
use_kafka = True
|
| 37 |
-
print("local kafka connection success!")
|
| 38 |
-
except BaseException:
|
| 39 |
-
print("local Kafka connection error")
|
| 40 |
-
print(f"Kafka status: {use_kafka}")
|
| 41 |
-
|
| 42 |
app = FastAPI()
|
| 43 |
|
| 44 |
PROJECT_ROOT = os.path.realpath(os.path.join(os.path.dirname(__file__), ".."))
|
|
|
|
| 6 |
import pandas as pd
|
| 7 |
from fastapi import FastAPI
|
| 8 |
from gamma.utils import association
|
|
|
|
| 9 |
from pydantic import BaseModel
|
| 10 |
from pyproj import Proj
|
| 11 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 12 |
app = FastAPI()
|
| 13 |
|
| 14 |
PROJECT_ROOT = os.path.realpath(os.path.join(os.path.dirname(__file__), ".."))
|