Friday, 28 April 2023

HUDI

 Question on Hudi 

What is Copy on write and Merge on Read 

What is Serialization in Athena table Serde 

  • serde tells us format of input file. If we have json file and we load into hive table. serde tells it the underlying file is in json format. 

What happens if we remove or add a column in Glue 

- Materialized views and joining tables 

-- how hudi knows something is latest record 

-- Hudi Snapshot and Incremental API 

How data skews are handled in Spark 

  • SALT Technique in spark for Handling Skews - Adding a new field to join key 

-- How do we partition table in S3 or HUDI , what partition keys do we use. What are indexes in HUDI 


How Parquet stores data -


Monday, 24 April 2023

Data Streaming

 There are various data streaming platforms 

- KAFKA

- Apark Streaming 

- Flink 

- Glue for Streaming 


Amazon Kinesis Data Streams 


Understanding AVRO 

- https://www.confluent.io/blog/avro-kafka-data/

- Avro has metadata in JSON and data is stored in Binary. Each file has header which is human readable and rest is binary 

-- Parquet and ORC - Are both suitable for Write once and Read heavy. ORC is optimized for hive with Hadoop. Both are columunar formas 

-- Parquet works best with Spark 

How Parquet Stores Data.Good Article  - https://www.linkedin.com/pulse/all-you-need-know-parquet-file-structure-depth-rohan-karanjawala/

- Both ORC and Parquet squeeze the data - ORC ( Optimized Row Columunar) 

https://www.upsolver.com/blog/the-file-format-fundamentals-of-big-data#:~:text=The%20ORC%20file%20format%20stores,reduce%20read%20and%20decompression%20loads.

How to Create Data streaming job 

https://aws.amazon.com/blogs/big-data/crafting-serverless-streaming-etl-jobs-with-aws-glue/

Blog on Stream Processin 

https://www.upsolver.com/wp/stream-processing-ebook?submissionGuid=d36d39a3-91ac-4f92-9a9b-f41cfd7eb305

Message Broker / Stream Processor 

- Apache Kafka 

- Kinesis Data streams 

Stream Processing Tools 

- AWS Kinesis 

- Apache Spark Streaming 

- Apache Flink 

- Kafka streaming API 

- Amazon Kinesis Data analytics allows you to process data in streams and use analytics over it. It uses apache flink

--------------------------------------------

For Loading data 

- Kinesis Firehose  -Firehose also has some tranformation capabilities . Firehose has dynamic partition capabilities and can partition data based on keys


Kafka Consumer API 

-- this is not streaming API 

-- this allows to read kafka streams and validate them or apply a logic 


Wednesday, 19 April 2023

Python Notes

 Quick python revision notes

# formating
name = 'bhagvant'
f = f"this is {name}"
print(f)

# Notice we are not adding f
name = 'bhagvant'
greeting = 'hello {}'
k = greeting.format(name)

################ LIST ###########################
l = ["Bob", "Rolf", "Anne"]

l[0] = "Smith"
l.append("Jen")

# extend allows to add set or tuple to list, l+s will not work , works only for list
s = {"Bob", "Rolf", "Anne"}
l.extend(s)
print(l)

# we can even insert a set to list
l[0] = s

#gets the length of list
len(l)

# Get number of times element appears in list
print(l.count("Bob"))

# to remove element from given inde type POP
l.pop(2)
l.pop() # removes last element

# reverse a list
l.reverse()

# sorts a list
prime_numbers = [11, 3, 7, 5, 2]
prime_numbers.sort()

############ Tuple ##############################

# you can concat 2 tuples but cant , change element of tuple
tuple1 = (0, 1, 2, 3)
tuple2 = ('python', 'geek')
 
# Concatenating above two
print(tuple1 + tuple2)



################# SET ##########################
# sets are unordered, Cannot contain duplicates and efficient for searching elements


friends = {"Bob", "Rolf", "Anne"}
abroad = {"Bob", "Anne"}
# to Create a set you cant use {} , this will create empty dictionary
a = set()

s.add("Jen")
# set cant have same element twice , its distinct
s.add("Bob")


print(friends.difference(abroad))
# returns empty set
print(abroad.difference(friends))

print(friends.intersection(abroad))


friends = {"Bob", "Rolf", "Anne"}
abroad = {"Bob", "Anne"}

# superset
if friends > abroad :
    print('superset')
   
# subset
if abroad < friends:
    print('subset')


# easy to check key in , fastest data structure due to hash table
for name in friends:
    if name in abroad:
        print(f' {name} gone abroad')



############# Is operator #######################
if 2 variables point to same object
x = 5
y = 5

print(x is y)

############ If condition #####################

dayofweek = 'Monday'

if dayofweek == 'Monday':
    print(1)

######### in keyword #################
# The `in` keyword works in most sequences like lists, tuples, and sets.

friends = ["Rolf", "Bob", "Jen"]
print("Jen" in friends)

# --

movies_watched = {"The Matrix", "Green Book", "Her"}
user_movie = input("Enter something you've watched recently: ")

print(user_movie in movies_watched)


############ LOOPS #######################

while n < 5:
    print(n)
    n+=1
   
while True:
    break
   
for i in range(10):
    print(f'for {i}')

#How to have 2 loops and loop through
for i in range(_size):
    k = i + 1
    for j in range(k, _size):

# remember when range start and end are equal it will not print
for i in range(1,1):
    print('wont print')

# incrementing range by 2 every time
for i in range(0,6,2):
    print(i)

########### List comprehension ###################

print([i for i in range(10) if i in [2,4,6]])


######## Dicationaries ##################

friend_ages = {"Rolf": 24, "Adam": 30, "Anne": 27}
print(friend_ages.keys())
print(friend_ages.values())
print(friend_ages['bhagvant'])

# adding to dictionary
friend_ages['bhagvant'] = 54
   
# Simple
student_attendance = {"Rolf": 96, "Bob": 80, "Anne": 100}
for student in student_attendance:
    print(f"{student}: {student_attendance[student]}")

# better
for i,k in friend_ages.items():
    print(i,k)


############## Functions #####################

def abc(a=1,b=2):
    print('this is function',a,b)
   
abc(1,3)

def abc(a=1,b=2):
    print('this is function',a,b)
    c = a+b
    return c
   
total = abc(1,3)



########## Lambda ############################
#map allows to add lambda to sequence
l = [1,2,3,4,5]
sum1 = lambda x:x+1

map_object = map(sum1,l)
   
print(list(map_object))



# * packs arguments into sinle list
a,*b = 1,2,3,4
print('first',a,b)

# unpacks it into tuple
def abc(k,*a):
    print(a)
    print(k)
   
abc(1,2,3,4)


# it packs the values into dictionary
def abc(**kwargs):
    print(kwargs)
   
abc(a='kk',b='jj')

anotherfunctionwithKeyValue(**kwargs)

def abc(**kwargs):
    print(kwargs)
    # allows to pass it
    anotherfunctionwithKeyValue(**kwargs)
   
abc(a='kk',b='jj')



########## Object oriented programming ###################

class abc:
    # notice init has to have self as argument
    def __init__(self,a=0,b=0):
        self.c = a+b
        # all class variables need self
   
    # self has to be passed as argument
    def multi(self):
        print(self.c)
     
    # used to print info about class , shoudl have a return
    def __str__(self):
        return f"value of a is {self.c}"
       
   
k = abc(1,1)
k.multi()
# str gets called when you print class ref
print(k)

Apache Airflow Notes

 Notes on apache airflow


with DAG( dag_id = ) as dag : 


- Operators in Airflow - default python , bash the default ones which come with it 

- https://airflow.apache.org/docs/apache-airflow/stable/_api/airflow/operators/index.html

- When interactin with any third party provides we install those providors 

- Example AWS , Snowflake , Databricks 

- Sensors - 

- Task group - Airflow utils has task group ids which can be used for grouping


Notes on Python

 My Notes on Python 

Interview Use cases to talk about

 Engineering Challenges solved 

Building Data lake 

- A system built over time was moved to Datalake 

- Initial system has multiple redshifts / Compaction issues / Multiple s3 paths from which data was consumed by customers

- Ownership was divided among teams but not logically 

- access was not controlled , redundant access 



STAR format - How was the impact measured 


- SLA improvement -- Team was able to redesign pipeline during migration to datalake and take our redundant steps improving sla 6 hours to 3 hours 

- Cost saving of 200k by moving processing


People Challenges solved 


Hiring and Recruiting Issues solved 


Big Projects Handled 


Appraisal Ratings 


Cost Saving 

- What is cost of Redshift 

- What is cost of EMR 

- Cost of Athena 

- Number of nodes 


Ra3.16x Large - Reserved instance - 75,000 - 48vcpu , 384 ram , 128TB space , Scales up to 16 petabytes 

DS2- 8x large - 16TB , HDD , 244gb memory ,  -- We had 50 Node cluster - DS2 is deprecated - 30 thousand 

DC2-8x - 32vcpu ,244 gb memory,  2.5TB SSD, 

EMR Type used - R5d - 48 VCPU, 512 memory - Supports upto EMR 6.3