· 8 years ago · Dec 05, 2017, 11:08 AM
1def gen_cassandra_createsql_by_df_schema(keyspace_name,table_name,df,override_fields={},primary_key_string=''):
2 """Generate cassandra create table sql string
3 based on DataFrame's table schema
4
5 :param keyspace_name: cassandra's keyspace name
6 :param table_name: cassandra's table name in keyspace_name
7 :param df: source dataframe
8 :param override_fields: key:value pairs for fields which type should be manually settled
9 key is column name, value is cassandra's datatype name
10 :param primary_key_string: primary key construct string
11 :return: CREATE TABLE sql string
12 """
13
14 known_types = {"StringType":'text','IntegerType':'int','BooleanType':'int'}
15 converted_fields ={}
16
17 schema = df.schema
18
19 #fill result_fields dict with original field names and types converted df -> cassandra
20 for field in schema.fields:
21 if field.name in override_fields:
22 converted_fields.update({field.name:override_fields[field.name]})
23 else:
24 dtype = str(field.dataType)
25 if dtype in known_types:
26 converted_fields.update({field.name:known_types[dtype]})
27 else:
28 raise ValueError("Unknown field.dataType \"{}\", add this type rule to the known_types dict".format(field.dataType))
29
30 #fill CREATE_TABLE string with fields definitions
31
32 db_fields = "".join(["%s %s, " % (k, v) for k, v in converted_fields.items()])
33 create_sql_statement = "CREATE TABLE IF NOT EXISTS {}.{} ( {}{} );".format(keyspace_name,table_name,db_fields,primary_key_string)
34 return create_sql_statement
35
36#example usage, trans_spark - my dataframe, kkbox/transactions - my keyspace and table name
37print(gen_cassandra_createsql_by_df_schema("kkbox","transactions"\
38 ,trans_spark,{'transaction_date':'date',"membership_expire_date":'date'}
39 ,"PRIMARY KEY (msno, transaction_date)",))
40#example output:
41CREATE TABLE IF NOT EXISTS kkbox.transactions ( msno text, payment_plan_days int, plan_list_price int, actual_amount_paid int,
42is_auto_renew int, transaction_date date, membership_expire_date date, is_cancel int, payment_method_id_1 int, payment_method_id_2 int,
43payment_method_id_3 int, payment_method_id_4 int, payment_method_id_5 int, payment_method_id_6 int, payment_method_id_7 int,
44payment_method_id_8 int, payment_method_id_10 int, payment_method_id_11 int, payment_method_id_12 int, payment_method_id_13 int,
45payment_method_id_14 int, payment_method_id_15 int, payment_method_id_16 int, payment_method_id_17 int, payment_method_id_18 int,
46payment_method_id_19 int, payment_method_id_20 int, payment_method_id_21 int, payment_method_id_22 int, payment_method_id_23 int,
47payment_method_id_24 int, payment_method_id_25 int, payment_method_id_26 int, payment_method_id_27 int, payment_method_id_28 int,
48payment_method_id_29 int, payment_method_id_30 int, payment_method_id_31 int, payment_method_id_32 int, payment_method_id_33 int,
49payment_method_id_34 int, payment_method_id_35 int, payment_method_id_36 int, payment_method_id_37 int, payment_method_id_38 int,
50payment_method_id_39 int, payment_method_id_40 int, payment_method_id_41 int, PRIMARY KEY (msno, transaction_date) );