I have datasets with identical schemas stored in folders that denote an id for the dataset, e.g.:
\11111\dataset
\11112\dataset
Where the '11111' etc. indicates the dataset id. I am trying to write a transform in code repository to loop through the datasets and append them all together. The following code works for this:
def create_outputs(dataset_ids):
transforms = []
for id in dataset_ids:
@transform_df(
Output(output_path + "/appended_dataset"),
input_path=Input(input_path + id + "/dataset"),
)
def compute(input_path):
return input_path
transforms.append(compute)
return transforms
id_list = ['11111','11112']
TRANSFORMS = create_outputs(id_list)
However, rather than having the id's hardcoded in the id_list, I would like to have a separate dataset that holds the dataset id's that need to be appended. I am having difficulty getting something that works.
I have tried the following code, where the id_list_dataset holds the ids to be included in the append:
# input dataset
id_list_dataset = ["ri.foundry.main.dataset.abcdefg"]
schema = T.StructType([
T.StructField('ID', T.StringType())
])
sc = SparkContext.getOrCreate()
rdd = sc.parallelize(id_list_dataset)
sqlContext = SQLContext(sc)
# define dataframe
temp_df = sqlContext.createDataFrame(rdd, schema)
# get list of ID's
id_list = temp_df.select('ID').collect
TRANSFORMS = create_outputs(id_list)
However, this is giving the following error:
TypeError: 'method' object is not iterable